如何使用PySpark将含Null字段的数据流写入Kafka并保留Null字段
解决PySpark写入Kafka时保留Null字段的问题
问题根源在于PySpark的to_json函数默认会忽略值为Null的字段,导致生成的JSON消息中看不到这些字段。要保留Null字段,只需修改to_json的调用参数,添加ignoreNulls=false配置。
修改后的核心代码
将原selectExpr中的to_json部分调整为:
data_frame.selectExpr( "CAST(id AS STRING) AS key", "to_json(struct(metadata,payload), map('ignoreNulls', 'false')) AS value" )
或者也可以用字符串形式传递配置:
data_frame.selectExpr( "CAST(id AS STRING) AS key", "to_json(struct(metadata,payload), 'ignoreNulls=false') AS value" )
效果说明
设置ignoreNulls=false后,to_json会保留所有包含Null值的字段,最终生成的JSON消息会包含类似"test": null的结构,符合需求。
内容的提问来源于stack exchange,提问作者Smaillns
相关产品推荐
相关产品推荐

