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

通过Kafka、JDBC源连接器与PySpark读取Postgres数据格式异常及Avro解析报错

问题描述

我在Postgres中创建了sample_a表并插入数据,配置JDBC源连接器(使用Avro转换器)将数据推送至Kafka,在Kafka Control Center中可正常查看数据,但通过PySpark读取时,将value转为字符串显示乱码。尝试使用Schema Registry获取的schema调用from_avro方法解析时,出现数组越界等解析错误,请问如何正确访问数据属性并解决该问题?

相关代码与配置

Postgres表创建语句

CREATE TABLE IF NOT EXISTS public.sample_a
(
    id text COLLATE pg_catalog."default" NOT NULL,
    is_active boolean NOT NULL,
    is_deleted boolean NOT NULL,
    created_by integer NOT NULL,
    created_at timestamp with time zone NOT NULL,
    created_ip character varying(30) COLLATE pg_catalog."default" NOT NULL,
    created_dept_id integer NOT NULL,
    updated_by integer,
    updated_at timestamp with time zone,
    updated_ip character varying(30) COLLATE pg_catalog."default",
    updated_dept_id integer,
    deleted_by integer,
    deleted_at timestamp with time zone,
    deleted_ip character varying(30) COLLATE pg_catalog."default",
    deleted_dept_id integer,
    sql_id bigint NOT NULL,
    ipa_no character varying(30) COLLATE pg_catalog."default" NOT NULL,
    pe_id bigint NOT NULL,
    uid character varying(30) COLLATE pg_catalog."default" NOT NULL,
    mr_no character varying(15) COLLATE pg_catalog."default" NOT NULL,
    site_id integer NOT NULL,
    entered_date date NOT NULL,
    CONSTRAINT pk_patient_dilation PRIMARY KEY (id)
);

数据插入语句

INSERT INTO sample_a (id, is_active, is_deleted, created_by, created_at, created_ip, created_dept_id, updated_by, updated_at, updated_ip, updated_dept_id, deleted_by, deleted_at, deleted_ip, deleted_dept_id, sql_id, ipa_no, pe_id, uid, mr_no, site_id, entered_date)
VALUES ('00037167-0894-4373-9a56-44c49d2285c9', TRUE, FALSE, 70516, '2024-10-05 08:12:25.069941+00','10.160.0.76', 4, 70516, '2024-10-05 09:25:55.218961+00', '10.84.0.1',4,NULL, NULL, NULL, NULL, 0,0,165587147,'22516767','P5942023',1,'10/5/24');

JDBC源连接器配置

{
  "name": "JdbcSourceConnectorConnector_0",
  "config": {
    "value.converter.schema.registry.url": "http://schema-registry:8081",
    "key.converter.schema.registry.url": "http://schema-registry:8081",
    "name": "JdbcSourceConnectorConnector_0",
    "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
    "key.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "connection.url": "jdbc:postgresql://postgres:5432/",
    "connection.user": "postgres",
    "connection.password": "********",
    "table.whitelist": "sample_a",
    "mode": "bulk"
  }
}

PySpark读取代码(乱码版本)

from pyspark.sql.session import SparkSession
from pyspark.sql.functions import col

spark = SparkSession \
    .builder \
    .appName("Kafka_Test") \
    .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0") \
    .getOrCreate()

df = spark \
    .readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe", "sample_a") \
    .option("startingOffsets","latest") \
    .load()

df.selectExpr("cast(value as string) as value").writeStream.format("console").start()

spark.streams.awaitAnyTermination()

读取输出(乱码)

H00037167-0894-4373-9a56-44c49d2285c9?ڹ??d10.160.0.7??????d10.84.0.0????22516767P5942023¸  

尝试Avro解析的代码及错误

jsonSchema = {"type":"record","name":"sample_a","fields":[{"name":"id","type":"string"},{"name":"is_active","type":"boolean"},{"name":"is_deleted","type":"boolean"},{"name":"created_by","type":"int"},{"name":"created_at","type":{"type":"long","connect.version":1,"connect.name":"org.apache.kafka.connect.data.Timestamp","logicalType":"timestamp-millis"}},{"name":"created_ip","type":"string"},{"name":"created_dept_id","type":"int"},{"name":"updated_by","type":["null","int"],"default":None},{"name":"updated_at","type":["null",{"type":"long","connect.version":1,"connect.name":"org.apache.kafka.connect.data.Timestamp","logicalType":"timestamp-millis"}],"default":None},{"name":"updated_ip","type":["null","string"],"default":None},{"name":"updated_dept_id","type":["null","int"],"default":None},{"name":"deleted_by","type":["null","int"],"default":None},{"name":"deleted_at","type":["null",{"type":"long","connect.version":1,"connect.name":"org.apache.kafka.connect.data.Timestamp","logicalType":"timestamp-millis"}],"default":None},{"name":"deleted_ip","type":["null","string"],"default":None},{"name":"deleted_dept_id","type":["null","int"],"default":None},{"name":"sql_id","type":"long"},{"name":"ipa_no","type":"string"},{"name":"pe_id","type":"long"},{"name":"uid","type":"string"},{"name":"mr_no","type":"string"},{"name":"site_id","type":"int"},{"name":"entered_date","type":{"type":"int","connect.version":1,"connect.name":"org.apache.kafka.connect.data.Date","logicalType":"date"}}],"connect.name": "sample_a"}
df.select(from_avro("value", json.dumps(jsonSchema)).alias("sample_a")) \
    .select("sample_a.*").writeStream.format("console").start()

错误信息

org.apache.spark.SparkException: Malformed records are detected in record parsing. Current parse Mode: FAILFAST. To process malformed records as null result, try setting the option 'mode' as 'PERMISSIVE'.
    at org.apache.spark.sql.avro.AvroDataToCatalyst.nullSafeEval(AvroDataToCatalyst.scala:113)
    at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown Source)
    at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
    at org.apache.spark.sql.execution.WholeStageCodegenExec$$anon$1.hasNext(WholeStageCodegenExec.scala:760)
    at org.apache.spark.sql.execution.datasources.v2.DataWritingSparkTask$.$anonfun$run$1(WriteToDataSourceV2Exec.scala:435)
    at org.apache.spark.util.Utils$.tryWithSafeFinallyAndFailureCallbacks(Utils.scala:1538)
    at org.apache.spark.sql.execution.datasources.v2.DataWritingSparkTask$.run(WriteToDataSourceV2Exec.scala:480)
    at org.apache.spark.sql.execution.datasources.v2.V2TableWriteExec.$anonfun$writeWithV2$2(WriteToDataSourceV2Exec.scala:381)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
    at org.apache.spark.scheduler.Task.run(Task.scala:136)
    at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:548)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1504)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:551)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
    at java.base/java.lang.Thread.run(Thread.java:829)
Caused by: java.lang.ArrayIndexOutOfBoundsException: Index 70516 out of bounds for length 2
    at org.apache.avro.io.parsing.Symbol$Alternative.getSymbol(Symbol.java:460)
    at org.apache.avro.io.ResolvingDecoder.readIndex(ResolvingDecoder.java:283)

解决方案

核心问题分析

  1. 乱码原因:Avro是二进制序列化格式,直接将其转为字符串必然出现乱码,属于正常现象,无需针对性解决,正确解析后即可获取明文数据。
  2. 解析错误原因:Confluent JDBC连接器写入的Avro消息采用Confluent Avro格式,包含1字节魔术位+4字节Schema ID的前缀;而Spark原生from_avro仅支持解析纯Avro二进制数据,无法识别该前缀,因此触发数组越界等解析错误。

具体解决步骤

1. 添加兼容依赖

在构建SparkSession时,需要添加Confluent Schema Registry客户端和Spark Avro的兼容依赖,版本需与Spark、Confluent平台版本匹配(示例使用Spark 3.3.0 + Confluent 7.3.0):

from pyspark.sql.session import SparkSession
from pyspark.sql.functions import col, from_avro
import json

spark = SparkSession \
    .builder \
    .appName("Kafka_Test") \
    .config("spark.jars.packages", 
            "org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0,"
            "io.confluent:kafka-avro-serializer:7.3.0,"
            "io.confluent:kafka-schema-registry-client:7.3.0,"
            "org.apache.spark:spark-avro_2.12:3.3.0") \
    .getOrCreate()

# 读取Kafka数据流
df = spark \
    .readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe", "sample_a") \
    .option("startingOffsets","latest") \
    .load()

2. 正确解析Confluent Avro数据

提供两种可行的解析方式:

方式一:通过Schema Registry自动拉取schema

直接指定Registry地址,让函数自动获取对应topic的schema并处理前缀:

# 配置Schema Registry地址
avro_options = {"schema.registry.url": "http://schema-registry:8081"}

# 解析value字段,指定topic的value schema标识(格式:{topic名称}-value)
parsed_df = df.select(
    from_avro(col("value"), "sample_a-value", avro_options).alias("sample_a")
)

# 展开所有字段并输出
parsed_df.select("sample_a.*").writeStream \
    .format("console") \
    .option("truncate", False) \
    .start()

spark.streams.awaitAnyTermination()

方式二:手动指定schema字符串(生产环境推荐)

若已从Schema Registry获取到schema的JSON字符串,可直接传入并配置Registry地址处理前缀:

# 从Schema Registry获取的schema JSON
jsonSchema = {"type":"record","name":"sample_a","fields":[{"name":"id","type":"string"},{"name":"is_active","type":"boolean"},{"name":"is_deleted","type":"boolean"},{"name":"created_by","type":"int"},{"name":"created_at","type":{"type":"long","connect.version":1,"connect.name":"org.apache.kafka.connect.data.Timestamp","logicalType":"timestamp-millis"}},{"name":"created_ip","type":"string"},{"name":"created_dept_id","type":"int"},{"name":"updated_by","type":["null","int"],"default":None},{"name":"updated_at","type":["null",{"type":"long","connect.version":1,"connect.name":"org.apache.kafka.connect.data.Timestamp","logicalType":"timestamp-millis"}],"default":None},{"name":"updated_ip","type":["null","string"],"default":None},{"name":"updated_dept_id","type":["null","int"],"default":None},{"name":"deleted_by","type":["null","int"],"default":None},{"name":"deleted_at","type":["null",{"type":"long","connect.version":1,"connect.name":"org.apache.kafka.connect.data.Timestamp","logicalType":"timestamp-millis"}],"default":None},{"name":"deleted_ip","type":["null","string"],"default":None},{"name":"deleted_dept_id","type":["null","int"],"default":None},{"name":"sql_id","type":"long"},{"name":"ipa_no","type":"string"},{"name":"pe_id","type":"long"},{"name":"uid","type":"string"},{"name":"mr_no","type":"string"},{"name":"site_id","type":"int"},{"name":"entered_date","type":{"type":"int","connect.version":1,"connect.name":"org.apache.kafka.connect.data.Date","logicalType":"date"}}],"connect.name": "sample_a"}

schema_str = json.dumps(jsonSchema)
avro_options = {"schema.registry.url": "http://schema-registry:8081"}

# 解析value字段
parsed_df = df.select(
    from_avro(col("value"), schema_str, avro_options).alias("sample_a")
)

# 展开所有字段并输出
parsed_df.select("sample_a.*").writeStream \
    .format("console") \
    .option("truncate", False) \
    .start()

spark.streams.awaitAnyTermination()

3. 异常数据处理(可选)

生产环境中可配置PERMISSIVE模式,将解析失败的记录转为null,避免任务中断:

from_avro(col("value"), schema_str, avro_options, mode="PERMISSIVE")

内容的提问来源于stack exchange,提问作者RushHour

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 03:04:51