【发布时间】: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