You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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解析即可:

  1. 准备对应的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"}
  ]
}
"""
  1. 修改代码中的解析逻辑:
# 替换原来的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类型:

  1. 定义目标Struct的Schema:
customer_schema = StructType([
    StructField("firstname", StringType(), True),
    StructField("ip", StringType(), True),
    StructField("id", LongType(), True),
    StructField("lastname", StringType(), True)
])
  1. 添加解析逻辑:
# 替换原来的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.25 17:01:36