通过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)
解决方案
核心问题分析
- 乱码原因:Avro是二进制序列化格式,直接将其转为字符串必然出现乱码,属于正常现象,无需针对性解决,正确解析后即可获取明文数据。
- 解析错误原因: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
相关产品推荐
相关产品推荐

