如何读取Struct类型DataFrame列并将键值对添加至Kafka Headers?
解决方案
你只需要先从Row中取出Headers这个Struct对应的子Row,再从这个子Row里读取具体的字段值即可,修改后的代码如下:
df.foreachPartition { partitionData: Iterator[Row] => partitionData.foreach { row => val kafkaProps = getKafkaConfig() val producer = JsonObjectProducer.getOrCreateProducer(kafkaProps) val data = new ProducerRecord[String, String](kafkaTopic, row.getAs[String]("Payload")) if(headerFlag == "True"){ // 先获取Struct类型的Headers对应的Row val headersRow = row.getAs[Row]("Headers") // 从子Row中读取各个字段并添加到Kafka Headers data.headers().add("Id", headersRow.getAs[String]("Id").getBytes(StandardCharsets.UTF_8)) data.headers().add("key", headersRow.getAs[String]("key").getBytes(StandardCharsets.UTF_8)) data.headers().add("ver", headersRow.getAs[String]("ver").getBytes(StandardCharsets.UTF_8)) } producer.send(data, new Callback() { def onCompletion(rm: RecordMetadata, e: Exception) { if (e != null) { e.printStackTrace() } } }) } }
关键说明
- 在Spark中,Struct类型的列会被封装成
Row对象,所以需要先用row.getAs[Row]("Headers")取出整个Struct结构 - 之后就可以像操作普通Row一样,用
getAs[String]从这个子Row里获取对应的字段值,再转成字节数组添加到Kafka Headers中
内容的提问来源于stack exchange,提问作者abhimanyu
相关产品推荐
相关产品推荐

