【问题标题】:Messages are not sent to Kafka from PySpark消息不会从 PySpark 发送到 Kafka
【发布时间】:2019-01-24 09:46:56
【问题描述】:

我想使用 Python Kafka 连接器将数据发送到 Kafka。当我从pyspark shell 运行代码时,一切正常。 但是,当我以spark-submit 运行它时,不会发送消息。日志中没有错误,程序执行显示为成功。但是消息不会发送到 Kafka。

import json
import datettime
from kafka import KafkaProducer

producer = KafkaProducer(bootstrap_servers='XXX.XX.XXX.XXX:9092')
end = datetime.datetime.now().isoformat()
country = "es"
message = {'country': country, 'end': end, 'status': '1'}
msg = json.dumps(message)
print(msg)
producer.send('testtopic', msg)

我不明白为什么会这样。下面我提供spark-submit的参数:

spark-submit \
--master yarn \
--deploy-mode cluster \
--driver-memory 11g \
--driver-cores 3 \
--num-executors 6 \
--executor-memory 6g \
--executor-cores 2 \
--conf spark.dynamicAllocation.enabled=false \
--conf spark.sql.broadcastTimeout=1500 \
--queue t1 \
s3://my-test-bucket/test1/test.py

【问题讨论】:

    标签: python python-3.x apache-spark pyspark apache-kafka


    【解决方案1】:

    我不得不在producer.send('testtopic', msg) 之后使用producer.flush()。只有在这种情况下,当我使用spark-submit 运行代码时,才会将消息发送到 Kafka 队列。 否则,消息不会发送。

    但是,奇怪的是,从 pyspark shell 执行代码时不需要producer.flush()

    【讨论】:

      【解决方案2】:

      生产者从批处理队列中轮询一批消息,每个分区一个批处理。当满足以下条件之一时,批次已准备就绪:

      batch.size 已达到。注意:较大的批次通常具有更好的压缩比和更高的吞吐量,但它们具有更高的延迟。

      linger.ms(基于时间的批处理阈值)已达到。注意:没有设置 linger.ms 值的简单指南;您应该测试特定用例的设置。对于小型事件(100 字节或更少),此设置似乎没有太大影响。

      同一代理的另一批已准备就绪。

      生产者调用flush()或close()。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2017-02-23
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2021-06-12
        • 1970-01-01
        相关资源
        最近更新 更多