【发布时间】:2019-11-30 04:49:02
【问题描述】:
我有一个任务,我想从 kafka 读取数据并使用 spark spark 流处理它,我想将数据发送到 Hbase。
在spark官方文档中,我发现:
def sendPartition(iter):
# ConnectionPool is a static, lazily initialized pool of connections
connection = ConnectionPool.getConnection()
for record in iter:
connection.send(record)
# return to the pool for future reuse
ConnectionPool.returnConnection(connection)
dstream.foreachRDD(lambda rdd: rdd.foreachPartition(sendPartition))
但我找不到任何线索来使用 pyspark 将 ConnectionPool 设置为 Hbase。
我也不明白流媒体是如何工作的?
在代码中有foreachPartition,我想清楚这些分区是否在同一个火花容器上?
为每个rdd的每个分区重置闭包中的所有变量?
有没有办法在工人级别设置变量?
剂量globals() 是工人级别吗?还是集群级别的?
【问题讨论】:
标签: apache-spark pyspark spark-streaming