Spark DataFrame写入Kafka遇序列化问题 求并行解决方案
为什么会出现序列化错误?
你遇到的PicklingError本质原因是KafkaProducer对象无法被序列化并传递到Worker节点。
你在Driver端创建了producer实例,然后在custom_fun里引用了它。Spark执行foreach时,需要把这个函数以及它依赖的所有变量(包括producer)序列化后发送到各个Worker节点执行。但KafkaProducer内部包含了很多不能被Python的pickle机制序列化的组件——比如底层的网络连接句柄、itertools.count类型的计数器(用来生成消息ID)等,这些资源是和Driver进程绑定的,无法被安全地序列化和传输到Worker。
几种可行的并行解决方案
方案1:在Worker端按需创建KafkaProducer(foreach版本)
不要在Driver端创建Producer,而是在Worker的任务函数内部延迟初始化Producer,这样每个Worker的任务进程会自己创建独立的Producer,不需要序列化Driver端的实例。为了避免每条数据都创建一个Producer(太浪费资源),可以用函数属性做单例缓存:
from kafka import KafkaProducer import util df = sqlContext.createDataFrame([("foo", 1), ("bar", 2), ("baz", 3)], ('k', 'v')) def send_row(row): # 每个任务进程只初始化一次Producer if not hasattr(send_row, 'producer'): send_row.producer = KafkaProducer(bootstrap_servers=util.get_broker_metadata()) # 注意要把字符串转成字节流,KafkaProducer.send需要bytes类型 send_row.producer.send('topic', str(row.asDict()).encode('utf-8')) # 不要每条数据都flush,会严重影响性能,让Producer自动处理批量发送 # send_row.producer.flush() df.foreach(send_row)
方案2:用foreachPartition优化性能(推荐)
foreach是针对每条数据执行一次函数,而foreachPartition是针对整个数据分区执行一次函数。我们可以在每个分区的开头创建一个Producer,处理完整个分区的所有数据后再flush并关闭,这样能大幅减少Producer的创建开销,提升写入效率:
from kafka import KafkaProducer import util df = sqlContext.createDataFrame([("foo", 1), ("bar", 2), ("baz", 3)], ('k', 'v')) def send_partition(rows): producer = KafkaProducer(bootstrap_servers=util.get_broker_metadata()) for row in rows: producer.send('topic', str(row.asDict()).encode('utf-8')) # 整个分区数据处理完再flush,保证数据发送完成 producer.flush() producer.close() df.foreachPartition(send_partition)
方案3:使用Spark官方Kafka连接器(最优解)
Spark提供了原生的Kafka读写连接器(spark-sql-kafka-0-10模块),这是官方推荐的方式,它已经帮你处理了Producer的生命周期、批量写入、容错等问题,扩展性和性能都是最好的:
from pyspark.sql.functions import to_json, struct df = sqlContext.createDataFrame([("foo", 1), ("bar", 2), ("baz", 3)], ('k', 'v')) # 将DataFrame转换为Kafka要求的格式:必须包含value列(可选key列),且值为字符串或字节 kafka_ready_df = df.select( to_json(struct("k", "v")).alias("value") # 把行数据转成JSON字符串作为value ) # 写入Kafka kafka_ready_df.write \ .format("kafka") \ .option("kafka.bootstrap.servers", util.get_broker_metadata()) \ .option("topic", "topic") \ .save()
方案对比
- 方案1适合简单场景,但性能不如方案2;
- 方案2手动管理Producer,灵活性高,性能较好;
- 方案3是最稳定、高效的选择,尤其是在大数据量场景下,推荐优先使用。
内容的提问来源于stack exchange,提问作者Nachiket Kate

