如何按key提取Kafka header值作为Spark Dataset单独列
Spark Structured Streaming 按Key提取Kafka Header实现
配置说明
对接开启Header功能的Kafka数据流时,需要添加如下配置项,让Spark将Header信息加载到Dataset中:
.option("includeHeaders", true)
启用该配置后,headers会作为Dataset的顶层列存储,列内容为包含key、value字段的结构体数组,对应Schema结构如下:
root |-- topic: string (nullable = true) |-- key: string (nullable = true) |-- value: string (nullable = true) |-- timestamp: string (nullable = true) |-- headers: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- key: string (nullable = true) | | |-- value: binary (nullable = true)
原有方案缺陷
初期实现时通过数组下标定位目标header,代码如下:
val controlDataFrame = spark .readStream .format("kafka") .option("kafka.bootstrap.servers", kafkaLocation) .option("includeHeaders", true) .option("failOnDataLoss", value = false) .option("subscribe", "mytopic") .load() .withColumn("acceptTimestamp", element_at(col("headers"),1)) .withColumn("acceptTimestamp2", col("acceptTimestamp.value").cast("STRING"))
该方案健壮性不足:生产端版本更新时可能调整headers的排列顺序,只有header的key名称是稳定不变的,依赖下标定位极易出现字段提取错误。
优化实现
通过map_from_entries函数将存储headers的结构体数组转换为Map结构,即可直接通过稳定的key名称提取目标header,无需依赖数组顺序,核心实现代码如下:
.withColumn("headers1", map_from_entries(col("headers"))) .withColumn("acceptTimestamp2", col("headers1.acceptTimestamp").cast("STRING"))
内容的提问来源于stack exchange,提问作者Eljah
相关产品推荐
相关产品推荐

