news 2026/7/28 11:35:55

基于PySpark整合Spark Streaming与Kafka

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
基于PySpark整合Spark Streaming与Kafka

本文内容主要给出基于PySpark程序,整合Spark Streaming和Kafka,实现实时消费和处理topic消息,为PySpark开发大数据实时计算项目提供基本参考。(未来将陆续更新基于Scala开发大数据实时计算项目的文章)

1 程序环境准备:

这里不再使用Spark的集群环境,因涉及的计算资源测试环境受限,目前两台虚拟机:1个vcore+2G内存,其中一台虚拟机启动Spark Streaming服务进程,另外一台虚拟机启动kafka进程。
虚拟机A:启动单实例kafka服务
虚拟机B:运行PySpark程序
在VM A,程序环境要求安装jdk1.8以上以及与kafka匹配版本的scala版本
版本兼容说明:

kafka:kafka_2.11-2.4.0 java:java version "1.8.0_11" scala: Scala 2.12.0

这里需要注意:如果使用kafka_2.12版本以上,需要使用jdk1.8.0_212以上;kafka_2.12与jdk1.8.0_11有不兼容地方,kafka启动报错提示java.lang.VerifyError: Uninitialized object exists on backward branch 209

1.1 基本配置

(1)配置单机zk这里无需依赖ZooKeeper集群,只需使用kafka自带的zk服务即可
vim /opt/kafka_2.11-2.4.0/config/zookeeper.properties

dataDir=/opt/zookeeper # zk的snapshot数据存储路径 clientPort=2181 # 按默认端口

(2)配置kafka的,路径/opt/kafka_2.11-2.4.0/config/ server.properties

log.dirs=/opt/kafka-logs # 存放kafka数据目录 zookeeper.connect=127.0.0.1:2181 # 按默认连接本机zk即可
1.2 启动zk和kafka
[root@nn kafka_2.11-2.4.0]# pwd/opt/kafka_2.12-2.4.0[root@nn kafka_2.11-2.4.0]# nohup ./bin/zookeeper-server-start.sh config/zookeeper.properties 2>&1 &

kafka server后台启动:

[root@nn kafka_2.11-2.4.0]# nohup bin/kafka-server-start.sh config/server.properties 2>&1 &
1.3 测试单实例Kafka

对于kafka单节点而言,这里只能使用1个分区且1个replication-factor,topic名称为sparkapp

[root@nn kafka_2.11-2.4.0]# ./bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic sparkappCreated topic sparkapp.

打开一个新的shell,用于启动producer

[root@nn kafka_2.11-2.4.0]# bin/kafka-console-producer.sh --broker-list localhost:9092 --topic sparkapp

再打开一个新的shell,用于启动consumer

[root@nn kafka_2.11-2.4.0]# bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic sparkapp

在producer shell输入字符串,consumer端可以看到相应输出,说明单机的kafka可以正常运行,下面将使用Spark Streaming实时读取kafka的输入流

2 整合streaming和kafka
2.1 配置依赖包

具体说明参考官方文档spark streaming连接kafka需要依赖两个jar包(注意版本号):
spark-streaming-kafka-0-8-assembly_2.11-2.4.3.jar: 下载链接
spark-streaming-kafka-0-8_2.11-2.4.4.jar: 下载链接
将这两个jar包放在spark 的jars目录下,需要注意的是:这两个jar包缺一不可,如果是在Spark集群上做测试,那么每个Spark节点都需要放置这两个jars包:

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/7/28 11:35:53

终极指南:如何用Applite彻底改变你的Mac软件管理体验

终极指南:如何用Applite彻底改变你的Mac软件管理体验 【免费下载链接】Applite User-friendly GUI macOS application for Homebrew Casks 项目地址: https://gitcode.com/gh_mirrors/ap/Applite 你是否厌倦了在Mac上手动下载、安装和管理软件的繁琐过程&…

作者头像 李华
网站建设 2026/7/28 11:35:19

C++ std::generate算法详解:从原理到实战的高效数据生成指南

1. 项目概述:为什么我们需要 std::generate ? 在C的日常开发中,尤其是处理容器初始化、数据填充或者生成测试数据时,我们经常会遇到一个看似简单但写起来有点啰嗦的场景:如何用一个特定的规则,去填充一个…

作者头像 李华
网站建设 2026/7/28 11:34:36

NBM5100A与PIC18F97J60在低功耗物联网设备中的协同设计

1. NBM5100A与PIC18F97J60的协同设计背景在低功耗物联网设备设计中,CR2032等纽扣电池面临着两大核心挑战:一是高内阻导致的脉冲负载能力不足,二是有限的容量难以满足长期工作需求。Nexperia推出的NBM5100A电池增强器与Microchip的PIC18F97J60…

作者头像 李华
网站建设 2026/7/28 11:32:49

中文NLP优势与AI开发实战指南

1. 汉字与AI的千年之约 汉字作为世界上唯一持续使用至今的古老文字系统,其独特结构为现代AI发展提供了天然优势。与拼音文字不同,每个汉字都是独立的信息单元,包含形、音、义三重特征。这种高密度信息载体特性,使得中文在自然语言…

作者头像 李华
网站建设 2026/7/28 11:32:46

物联网设备初级电池寿命优化方案与实践

1. 不可充电初级电池的寿命挑战与解决思路 在物联网设备和便携式电子设备中,不可充电的初级电池(如锂亚硫酰氯电池、碱性电池等)仍然是许多场景下的首选电源方案。这类电池具有能量密度高、自放电率低、工作温度范围广等优势,但也…

作者头像 李华
网站建设 2026/7/28 11:31:45

NBM7100A与MKV42F256VLH16在低功耗物联网设备中的协同优化

1. 硬币电池寿命延长器的核心挑战在低功耗物联网设备中,CR2032这类不可充电的硬币电池面临着两大核心挑战:高脉冲电流需求导致的电压骤降和有限容量下的快速耗尽。当无线模块启动传输或传感器进行数据采集时,瞬间电流可能达到电池标称容量的数…

作者头像 李华