如何实现PySpark Streaming应用:CSV转Kafka并行生产及热门姓名统计
PySpark Streaming 跨数据集热门姓名统计实现方案
一、Kafka生产者:并行读取CSV并发送至对应主题
首先安装依赖:
pip install kafka-python
优化后的生产者代码(修复序列化问题,确保消息发送完成):
from kafka import KafkaProducer import csv import json from concurrent.futures import ThreadPoolExecutor # 初始化Kafka生产者,用JSON序列化数据 producer = KafkaProducer( bootstrap_servers=["localhost:9092"], value_serializer=lambda v: json.dumps(v).encode('utf-8') ) def publish_to_topic(topic, data): """将单条数据发送至指定Kafka主题""" producer.send(topic, value=data) def process_csv(file_path, topic, name_col): """读取CSV文件并发送数据到对应主题""" with open(file_path, newline="", encoding='utf-8') as csvfile: reader = csv.DictReader(csvfile) for row in reader: # 提取姓名字段,统一字段名 publish_data = {"name": row[name_col].strip()} publish_to_topic(topic, publish_data) print(f"文件 {file_path} 数据已全部发送至主题 {topic}") # 并行处理两个CSV文件 if __name__ == "__main__": with ThreadPoolExecutor(max_workers=2) as executor: # 处理A.csv,姓名列是name,发送到topic_a executor.submit(process_csv, "A.csv", "topic_a", "name") # 处理B.csv,姓名列是name_,发送到topic_b executor.submit(process_csv, "B.csv", "topic_b", "name_") # 等待所有消息发送完成并关闭生产者 producer.flush() producer.close()
关键说明
- 用
value_serializer统一做JSON序列化,避免手动转字符串的格式问题 - 封装
process_csv函数,复用CSV处理逻辑,统一姓名字段为name,方便后续Spark解析 - 最后调用
flush()确保所有消息发送到Kafka,再关闭生产者
二、PySpark Streaming:消费Kafka数据并统计热门姓名
依赖准备
确保PySpark环境已安装,运行时需要指定Kafka连接器包:
spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.0 your_script.py
(注意包版本要和你的Spark版本匹配)
实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col, count, from_json from pyspark.sql.types import StructType, StringType from pyspark.sql.functions import sum as spark_sum if __name__ == "__main__": # 初始化SparkSession spark = SparkSession.builder \ .appName("HotNameStreaming") \ .getOrCreate() spark.sparkContext.setLogLevel("WARN") # 定义数据Schema name_schema = StructType().add("name", StringType()) # 从Kafka读取两个主题的数据 kafka_df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "localhost:9092") \ .option("subscribe", "topic_a,topic_b") \ .load() # 解析Kafka消息中的JSON数据 parsed_df = kafka_df.select( from_json(col("value").cast(StringType()), name_schema).alias("data") ).select("data.name").filter(col("name").isNotNull()) # 按姓名分组统计累计出现次数 total_count_df = parsed_df.groupBy("name") \ .agg(count("name").alias("total_count")) \ .orderBy(col("total_count").desc()) # 输出结果到控制台(也可以替换为数据库/文件输出) query = total_count_df.writeStream \ .outputMode("complete") \ .format("console") \ .option("truncate", False) \ .start() # 等待流处理结束 query.awaitTermination()
关键说明
- 使用
readStream读取多个Kafka主题,用逗号分隔主题名 - 定义Schema解析JSON数据,确保类型安全
- 通过
groupBy和agg统计姓名出现次数,按热度降序排列 outputMode("complete")会输出所有姓名的累计计数,适合全局统计场景- 可根据需求替换输出方式,比如写入Hive、MySQL或分布式文件系统
内容的提问来源于stack exchange,提问作者Murat
相关产品推荐
相关产品推荐

