【问题标题】:Way to use previous data with current data in pyspark with kafka stream使用kafka流在pyspark中使用先前数据和当前数据的方法
【发布时间】:2019-11-23 06:59:06
【问题描述】:

我正在从我的生产者发送 dict 对象并使用 pyspark 创建一个新对象。但是我想要形成的那种 obj 也需要以前数据的键值对。我尝试了窗口批处理和 reduceByKey,但它们似乎都不起作用。

假设我的生产者对象就像“url_id”和“url”对的列表。例如 {"url_id": "google.com"} 和 spark 我想形成一个对象,如: {"data": {"url_id": "url", "url_id_of_previous_url": "url",... .等等}

我的火花代码是:

conf = SparkConf().setAppName(appName).setMaster("local[*]")
        sc = SparkContext(conf=conf)

        stream_context = StreamingContext(sparkContext=sc, batchDuration=batchTime)
        kafka_stream = KafkaUtils.createDirectStream(ssc=stream_context, topics=[topic], 
                                          kafkaParams={"metadata.broker.list":"localhost:9092", 
                                                     'auto.offset.reset':'smallest'})
        lines = kafka_stream.map(lambda x: json.loads(x[1]))

在这之后我被困住了。你能告诉我用火花形成这样的obj是否可能吗?如果是,那我应该使用什么?

【问题讨论】:

    标签: apache-spark pyspark spark-structured-streaming


    【解决方案1】:

    据我所知,你可以通过两种方式解决这个问题,

    第一种方法很简单,让消息生成应用程序本身通过启用一些内部缓存来发送一对消息(当前和以前的)。

    第二种方法是使用 Spark Stateful Streaming 在 Spark 状态上下文中维护最后一条消息的值。当您使用 PySpark 时,我知道的唯一选择是使用 updateStateByKey 并启用检查点。

    PySpark Streaming 的典型流程如下,

    • 定义初始值和更新函数
    • 维护一个公共密钥来匹配当前和以前的消息,我在这个例子中使用了pair_msgs

      # RDD with initial state (key, value) pairs
      initialStateRDD = sc.parallelize([(u'pair_msgs', '{"url_id":"none"}')])
      
      def updateFunc(new_url_msg, last_url_msg):
          if not new_url_msg:
              return last_url_msg
          else:
              new_url_dict = json.loads(new_url_msg[0])
              new_url_dict['url_id_previous'] = json.loads(last_url_msg)['url_id']
              return json.dumps(new_url_msg)
      
    • 用公共键映射输入消息,在本例中为pair_msgs

    • 使用上述更新函数调用updateStateByKey转换。

      feeds = kafka_stream.map(lambda x: x[1])
      
      pair_feed = feeds.map(lambda feed_str: ('pair_msgs', feed_str)) \
                   .updateStateByKey(updateFunc, initialRDD=initialStateRDD)
      

    [注意:据我所知,PySpark Structured Streaming 还没有获得 Stateful Streaming 支持,所以我相信上面的例子仍然有意义]

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-10-30
      • 2021-02-09
      • 2018-06-21
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多