基于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
相关产品推荐
相关产品推荐

