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

基于PySpark Structured Streaming实现Kafka数据写入Delta Table及查询

PySpark结构化流读取Kafka写入Delta Table实现

完整代码实现

# 导入必要模块
from pyspark.sql.functions import col, current_timestamp, expr, explode, split, first

# 从Kafka主题ABC读取流数据
kafka_stream = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe", "ABC") \
    .option("startingOffsets", "earliest")  # 读取主题中所有历史数据
    .load()

# 处理数据:提取原始消息、解析字段、转换格式
processed_stream = kafka_stream \
    # 保留原始消息内容
    .withColumn("data", col("value").cast("string")) \
    # 将query string格式的消息拆分为键值对
    .withColumn("key_value", split(col("data"), "&")) \
    .withColumn("key_value_pair", explode(col("key_value"))) \
    .withColumn("key", split(col("key_value_pair"), "=")[0]) \
    .withColumn("value", split(col("key_value_pair"), "=")[1]) \
    # 按原始消息和Kafka时间戳分组,将键值对转为结构化字段
    .groupBy("data", col("timestamp").alias("kafka_timestamp")) \
    .pivot("key") \
    .agg(first("value")) \
    # 将company_name转换为StackOverflow格式
    .withColumn("company_name", expr("concat('Stack', upper(substring(company_name, 6, 1)), substring(company_name, 7))")) \
    # 生成处理时间戳(若需使用Kafka消息时间戳,替换为col("kafka_timestamp"))
    .withColumn("timestamp", current_timestamp()) \
    # 选择最终需要的字段
    .select("timestamp", "company_name", "data")

# 写入Delta Table
delta_write_query = processed_stream.writeStream \
    .format("delta") \
    .option("checkpointLocation", "/dbfs/tmp/kafka_to_delta_checkpoint")  # Databricks建议用DBFS路径存储checkpoint
    .table("kafka_abc_delta")  # 若表不存在会自动创建

# 启动流作业
delta_write_query.awaitTermination()

字段说明

  • timestamp:默认使用数据处理时的系统当前时间,若需使用Kafka消息自带的时间戳,将代码中.withColumn("timestamp", current_timestamp())替换为.withColumn("timestamp", col("kafka_timestamp"))即可。
  • company_name:通过字符串拼接逻辑将输入的stackoverflow转换为StackOverflow,如果有其他类似格式的公司名,可以调整expr中的表达式适配。
  • data:完整保留Kafka消息的原始字符串内容,确保原始数据可追溯。

控制台查看Delta Table数据的方法

方式1:直接查询Delta表

在Databricks Notebook中执行以下代码或SQL即可查看表中所有数据:

# PySpark查询
spark.read.format("delta").table("kafka_abc_delta").show()
-- SQL查询
SELECT * FROM kafka_abc_delta;

方式2:流作业实时控制台输出

如果需要在流作业运行时实时查看处理后的数据,可以在写入Delta的同时,新增一个控制台输出的流分支:

# 新增控制台输出流
console_write_query = processed_stream.writeStream \
    .format("console") \
    .outputMode("append") \
    .start()

# 等待两个流作业运行
delta_write_query.awaitTermination()
console_write_query.awaitTermination()

内容的提问来源于stack exchange,提问作者boring-coder

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 14:30:52