【问题标题】:pyspark streaming how to set ConnectionPoolpyspark流如何设置ConnectionPool
【发布时间】: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


    【解决方案1】:

    你需要使用 thrift 从 pyspark 与 Hbase 交互

    这里是一些代码参考

    http://shzhangji.com/blog/2018/04/22/connect-hbase-with-python-and-thrift/

    关于并行化,只需调用将在 foreach 或 foreachPartition 中发布到 Habse 的方法(在所有执行程序上以分布式方式执行),只需确保您的应用程序的每个任务/核心都有专用连接。

    【讨论】:

    • 我知道如何与Hbase交互,我的问题是如何通过所有Rdds在steaming中设置一个句柄。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-10-31
    • 2020-06-27
    • 1970-01-01
    • 2015-05-04
    • 1970-01-01
    • 2015-12-27
    • 1970-01-01
    相关资源
    最近更新 更多