Structured Streaming中如何从Binary列提取Struct字段?
问题:如何将Kafka流中的Binary类型value列转换为Struct类型?
我希望从Kafka接收流数据,并从value列中提取一个Struct字段,但无法将Binary列转换为Struct类型。请问如何将该value列转换为Struct类型?
我的脚本
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job from pyspark.sql.types import * from pyspark.sql.functions import * from pyspark.sql.avro.functions import from_avro ## @params: [JOB_NAME] args = getResolvedOptions(sys.argv, ['JOB_NAME']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "hoge:9092,fuga:9092") \ .option("subscribe", "customer") \ .option("startingOffsets", "earliest") \ .option("includeHeaders", "true") \ .load() df = df.selectExpr("CAST(value AS STRING)") query = df \ .writeStream \ .outputMode("append") \ .format("console") \ .option("truncate", "false") \ .option("checkpointLocation", "s3://bucketname/tmp/checkpoint/") \ .start() query.awaitTermination()
当前输出
+----------------------------------------------------------------------------------------------------------------------------------------------------------------+ |value | +----------------------------------------------------------------------------------------------------------------------------------------------------------------+ |Struct{firstname=Earl,ip=217.54.184.114,id=9425651,lastname=Haley} +----------------------------------------------------------------------------------------------------------------------------------------------------------------+
解决方案
根据你的输出格式,分两种场景给出处理方案:
场景1:Kafka消息是Avro序列化格式(推荐)
从你导入from_avro函数来看,大概率Kafka的value是Avro序列化的二进制数据,无需转成字符串,直接用from_avro解析即可:
- 准备对应的Avro Schema(替换为你实际的Schema):
avro_schema = """ { "type": "record", "name": "Customer", "fields": [ {"name": "firstname", "type": "string"}, {"name": "ip", "type": "string"}, {"name": "id", "type": "long"}, {"name": "lastname", "type": "string"} ] } """
- 修改代码中的解析逻辑:
# 替换原来的df = df.selectExpr("CAST(value AS STRING)") df = df.select(from_avro(col("value"), avro_schema).alias("customer")) # 提取Struct中的字段(可选) df = df.select("customer.firstname", "customer.ip", "customer.id", "customer.lastname")
场景2:Kafka消息是字符串格式的Struct
如果消息确实是Struct{key=value,...}这种字符串格式,需要通过字符串处理解析为Struct类型:
- 定义目标Struct的Schema:
customer_schema = StructType([ StructField("firstname", StringType(), True), StructField("ip", StringType(), True), StructField("id", LongType(), True), StructField("lastname", StringType(), True) ])
- 添加解析逻辑:
# 替换原来的df = df.selectExpr("CAST(value AS STRING)") df = df.select( # 去除首尾的Struct{} regexp_replace(col("value"), "^Struct\\{|\\}$", "").alias("cleaned_value") ).select( # 分割键值对并转换为Map,再强转为Struct map_from_entries( split(col("cleaned_value"), ",\\s*").map(lambda x: split(x, "=")) ).cast(customer_schema).alias("customer") ) # 提取Struct中的字段(可选) df = df.select("customer.*")
内容的提问来源于stack exchange,提问作者cal0708
相关产品推荐
相关产品推荐

