【发布时间】:2016-12-25 19:46:12
【问题描述】:
我有一个带有一些输入主题的 kafka 流。 这是我为接受 kafka 流而编写的代码。
conf = SparkConf().setAppName(appname)
sc = SparkContext(conf=conf)
ssc = StreamingContext(sc)
kvs = KafkaUtils.createDirectStream(ssc, topics,\
{"metadata.broker.list": brokers})
然后我创建两个原始流的键和值的 DStream。
keys = kvs.map(lambda x: x[0].split(" "))
values = kvs.map(lambda x: x[1].split(" "))
然后我在值 DStream 中执行一些计算。 例如,
val = values.flatMap(lambda x: x*2)
现在,我需要将键和 val DStream 组合起来,并以 Kafka 流的形式返回结果。
如何将val与对应的key结合起来?
【问题讨论】:
标签: python apache-kafka pyspark kafka-python