如何在PySpark中合并两个DStream?现有代码无输出问题求助
解决PySpark Streaming合并两个DStream无输出的问题
看起来你遇到的核心问题有两个:一是DStream的结构不符合join的要求,二是流数据可能不同步到达导致匹配失败,另外你需要的是合并JSON字段而不是简单的join结果元组。咱们一步步来修正:
问题分析
- 你当前的
dstream和dstream1只提取了id的值,变成了单值的DStream,但join操作要求每个DStream的元素是**(key, value)**键值对结构,这样Spark才能根据key关联两个流的数据。 - 即使结构正确,普通的
join是内连接,只有当两个流在同一批次同时出现相同key的数据时才会输出,若数据到达时间有差,就会无输出。 - 你需要的是合并两个JSON的字段,而不是join后得到的
(key, (value1, value2))元组。
修正方案
步骤1:调整DStream为键值对结构
把每个流的元素处理成(id, 完整JSON对象)的形式,这样后续可以用id作为关联键:
import json from pyspark.streaming.kafka import KafkaUtils from pyspark import SparkContext from pyspark.streaming import StreamingContext from pyspark.sql import SparkSession sc = SparkContext(appName="Sparkstreaming") spark = SparkSession(sc) sc.setLogLevel("WARN") ssc = StreamingContext(sc, 3) # 处理第一个Kafka流:(key, 完整JSON) kafka_stream = KafkaUtils.createStream(ssc,"localhost:2181","consumer-be-borsa",{"Be_borsa":1}) dstream = kafka_stream.map(lambda k, v: (json.loads(v)['id'], json.loads(v))) # 处理第二个Kafka流:(key, 完整JSON) kafka_stream1 = KafkaUtils.createStream(ssc,"localhost:2181","consumer-be-teleborsa",{"Be_teleborsa":1}) dstream1 = kafka_stream1.map(lambda k, v: (json.loads(v)['id'], json.loads(v)))
步骤2:使用全外连接避免数据丢失
用fullOuterJoin替代普通join,这样即使其中一个流的key暂时没匹配到,也能保留数据,后续再处理合并:
# 全外连接两个流 streamJoined = dstream.fullOuterJoin(dstream1)
步骤3:合并JSON字段
自定义一个函数,把两个JSON对象合并成一个,处理其中一方为None的情况:
def merge_json(record): key, (json1, json2) = record merged = {} # 合并第一个JSON的字段 if json1: merged.update(json1) # 合并第二个JSON的字段(重复字段会覆盖第一个,符合你的需求) if json2: merged.update(json2) return (key, merged) # 应用合并函数 merged_stream = streamJoined.map(merge_json)
步骤4:输出并启动流
用pprint()输出,并且用awaitTermination()替代time.sleep,这样流会一直运行直到手动停止:
merged_stream.pprint() ssc.start() ssc.awaitTermination() # 替代time.sleep,更可靠 # ssc.stop() # 手动停止时调用
关键说明
- 为什么用fullOuterJoin:如果两个流的同一key数据不是在同一个3秒批次到达,普通内连接会丢失数据,全外连接能确保所有数据都被处理。
- 合并逻辑:
update()方法会把第二个JSON的字段合并进去,如果有重复字段(比如id、Date),第二个流的字段会覆盖第一个,这和你的期望结果一致。 - 流的停止方式:
awaitTermination()会让流持续运行,直到你发送停止信号(比如Ctrl+C),比time.sleep更适合生产环境。
这样调整后,你就能得到合并了所有字段的JSON输出,符合你的期望结果。
内容的提问来源于stack exchange,提问作者Luca Lazzati
相关产品推荐
相关产品推荐

