【发布时间】:2016-03-29 10:31:33
【问题描述】:
阅读official docs 后,我尝试在火花流中使用checkpoint 和getOrCreate。一些sn-ps:
def get_ssc():
sc = SparkContext("yarn-client")
ssc = StreamingContext(sc, 10) # calc every 10s
ks = KafkaUtils.createDirectStream(
ssc, ['lucky-track'], {"metadata.broker.list": KAFKA_BROKER})
process_data(ks)
ssc.checkpoint(CHECKPOINT_DIR)
return ssc
if __name__ == '__main__':
ssc = StreamingContext.getOrCreate(CHECKPOINT_DIR, get_ssc)
ssc.start()
ssc.awaitTermination()
代码可以正常用于恢复,但恢复的上下文始终适用于旧进程函数。这意味着即使我更改了 map/reduce 功能代码,它也根本不起作用。
到目前为止,spark(1.5.2) 仍然不支持 python 的任意偏移量。那么,我应该怎么做才能使其正常工作?
【问题讨论】:
标签: python apache-spark apache-kafka pyspark spark-streaming