【问题标题】:How to combine two DStreams(pyspark)?如何结合两个 DStreams(pyspark)?
【发布时间】: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


    【解决方案1】:

    您可以在 2 个 DStream 上使用 join 运算符来合并它们。 当您进行映射时,您实际上是在创建另一个流。因此,join 将帮助您将它们合并在一起。

    例如:

    Joined_Stream = keys.join(values).(any operation like map, flatmap...)
    

    【讨论】:

    • 我没有得到这部分(any operation like map, flatmap...),你能详细说明一下吗。
    • 我不明白你想要实际做的事情(我提供了合并 2 个 DStreams 的通用答案)。问题是,如果您对值进行平面映射,则无法将它们映射回键,因为该输出将是一个扁平列表....通过合并 2 个 Dstream,您可以创建 RDD,每个 RDD 都包含两个键的元素& 值,只是不会有一对一的映射......
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-11-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-11-17
    相关资源
    最近更新 更多