You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 21:45:01