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