使用Apache Beam Python SDK读取Kafka时遇RuntimeException问题求助
问题:Apache Beam读取SASL_SSL认证的Avro格式Kafka数据触发RuntimeException
环境信息
- 平台:Apple M1
- MiniConda版本:22.11.1
- Python版本:3.9.16
- Apache Beam SDK版本:2.44.0
代码片段
import apache_beam as beam from apache_beam.io.external.kafka import ReadFromKafka from apache_beam.io.external.kafka import WriteToKafka from apache_beam.options.pipeline_options import PipelineOptions # some variables definition conf = { 'auto.offset.reset': 'latest', 'basic.auth.credentials.source': 'SASL_INHERIT', 'bootstrap.servers': '{}:{}'.format(host, sasl_port), 'client.id': 'demo-py-client', 'enable.auto.commit': 'true', 'group.id': 'group_id', 'isolation.level': 'read_uncommitted', 'sasl.mechanism': 'PLAIN', 'sasl.jaas.config': "org.apache.kafka.common.security.plain.PlainLoginModule required username='{}' password='{}';".format(username, password), 'sasl.password': password, 'sasl.username': username, 'schema.registry.client.cache.capacity': '1000', 'schema.registry.url': 'https://{}:{}@{}:{}'.format(username, password, host, 29650), 'security.protocol': 'SASL_SSL', 'enable.ssl.certificate.verification': 'false', 'specific.avro.reader': 'false' } with beam.Pipeline(options=PipelineOptions()) as pipeline: ( pipeline | 'Read' >> ReadFromKafka( consumer_config=conf, topics=[topic], with_metadata=False, key_deserializer='io.confluent.kafka.serializers.KafkaAvroDeserializer', value_deserializer='io.confluent.kafka.serializers.KafkaAvroDeserializer', ) | 'Print' >> beam.Map(print) )
错误栈信息
{ "name": "RuntimeError", "message": "java.lang.RuntimeException: Failed to build transform beam:transform:org.apache.beam:kafka_read_without_metadata:v1 from spec urn: \"beam:transform:org.apache.beam:kafka_read_without_metadata:v1 org.apache.beam.sdk.expansion.service.ExpansionService$ExternalTransformRegistrarLoader$1.getTransform(ExpansionService.java:151) org.apache.beam.sdk.expansion.service.ExpansionService$TransformProvider.apply(ExpansionService.java:400) org.apache.beam.sdk.expansion.service.ExpansionService.expand(ExpansionService.java:526) org.apache.beam.sdk.expansion.service.ExpansionService.expand(ExpansionService.java:606) org.apache.beam.model.expansion.v1.ExpansionServiceGrpc$MethodHandlers.invoke(ExpansionServiceGrpc.java:305) org.apache.beam.vendor.grpc.v1p48p1.io.grpc.stub.ServerCalls$UnaryServerCallHandler$UnaryServerCallListener.onHalfClose(ServerCalls.java:182) org.apache.beam.vendor.grpc.v1p48p1.io.grpc.internal.ServerCallImpl$ServerStreamListenerImpl.halfClosed(ServerCallImpl.java:354) org.apache.beam.vendor.grpc.v1p48p1.io.grpc.internal.ServerImpl$JumpToApplicationThreadServerStreamListener$1HalfClosed.runInContext(ServerImpl.java:866) org.apache.beam.vendor.grpc.v1p48p1.io.grpc.internal.ContextRunnable.run(ContextRunnable.java:37) org.apache.beam.vendor.grpc.v1p48p1.io.grpc.internal.SerializingExecutor.run(SerializingExecutor.java:133) java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) java.base/java.lang.Thread.run(Thread.java:829) Caused by: java.lang.RuntimeException: Couldn't resolve coder for Deserializer: class io.confluent.kafka.serializers.KafkaAvroDeserializer org.apache.beam.sdk.io.kafka.KafkaIO$Read$Builder.resolveCoder(KafkaIO.java:819) org.apache.beam.sdk.io.kafka.KafkaIO$Read$Builder.setupExternalBuilder(KafkaIO.java:750) org.apache.beam.sdk.io.kafka.KafkaIO$TypedWithoutMetadata$Builder.buildExternal(KafkaIO.java:1686) org.apache.beam.sdk.io.kafka.KafkaIO$TypedWithoutMetadata$Builder.buildExternal(KafkaIO.java:1677) org.apache.beam.sdk.expansion.service.ExpansionService$ExternalTransformRegistrarLoader$1.getTransform(ExpansionService.java:145)\n\t... 12 more }
环境依赖列表
Package Version ------------------------------- ---------- apache-beam 2.44.0 appnope 0.1.2 asttokens 2.0.5 avro 1.11.1 backcall 0.2.0 cachetools 4.2.4 certifi 2022.12.7 charset-normalizer 3.0.1 cloudpickle 2.2.1 comm 0.1.2 confluent-kafka 2.0.2 crcmod 1.7 debugpy 1.5.1 decorator 5.1.1 dill 0.3.1.1 docopt 0.6.2 entrypoints 0.4 executing 0.8.3 facets-overview 1.0.0 fastavro 1.7.0 fasteners 0.18 google-api-core 2.11.0 google-apitools 0.5.31 google-auth 2.16.0 google-auth-httplib2 0.1.0 google-cloud-bigquery 3.4.2 google-cloud-bigquery-storage 2.13.2 google-cloud-bigtable 1.7.3 google-cloud-core 2.3.2 google-cloud-dataproc 3.1.1 google-cloud-datastore 1.15.5 google-cloud-dlp 3.11.1 google-cloud-language 1.3.2 google-cloud-pubsub 2.14.0 google-cloud-pubsublite 1.6.0 google-cloud-recommendations-ai 0.7.1 google-cloud-spanner 3.27.0 google-cloud-videointelligence 1.16.3 google-cloud-vision 3.3.1 google-crc32c 1.5.0 google-resumable-media 2.4.1 googleapis-common-protos 1.58.0 grpc-google-iam-v1 0.12.6 grpcio 1.51.1 grpcio-status 1.48.2 hdfs 2.7.0 httplib2 0.20.4 idna 3.4 ipykernel 6.19.2 ipython 8.7.0 ipywidgets 8.0.4 jedi 0.18.1 jupyter-client 6.1.12 jupyter_core 5.1.1 jupyterlab-widgets 3.0.5 matplotlib-inline 0.1.6 nest-asyncio 1.5.6 numpy 1.22.4 oauth2client 4.1.3 objsize 0.6.1 orjson 3.8.5 overrides 6.5.0 packaging 23.0 pandas 1.5.3 parso 0.8.3 pexpect 4.8.0 pickleshare 0.7.5 pip 22.3.1 platformdirs 2.5.2 prompt-toolkit 3.0.36 proto-plus 1.22.2 protobuf 3.20.3 psutil 5.9.0 ptyprocess 0.7.0 pure-eval 0.2.2 pyarrow 9.0.0 pyasn1 0.4.8 pyasn1-modules 0.2.8 pydot 1.4.2 Pygments 2.11.2 pymongo 3.13.0 pyparsing 3.0.9 python-dateutil 2.8.2 pytz 2022.7.1 pyzmq 23.2.0 regex 2022.10.31 requests 2.28.2 rsa 4.9 setuptools 65.6.3 six 1.16.0 sqlparse 0.4.3 stack-data 0.2.0 timeloop 1.0.2 tornado 6.2 traitlets 5.7.1 typing_extensions 4.4.0 urllib3 1.26.14 wcwidth 0.2.5 wheel 0.37.1 widgetsnbextension 4.0.5 zstandard 0.19.0
问题原因分析
核心错误是Couldn't resolve coder for Deserializer: class io.confluent.kafka.serializers.KafkaAvroDeserializer,具体原因包括:
- Java依赖缺失:Apache Beam的Kafka外部Transform依赖Confluent的Avro序列化库,但默认的Beam Java扩展服务未包含这些依赖,导致无法加载
KafkaAvroDeserializer类。 - Coder未显式指定:使用Java端的Avro反序列化器时,Beam无法自动推断对应的Coder,导致数据格式解析失败。
- M1平台兼容性:ARM架构下部分Java依赖存在适配问题,可能引发类加载异常。
解决建议
1. 补充Java依赖并启动自定义扩展服务
Beam Python SDK调用Kafka外部Transform依赖Java扩展服务,需手动添加Confluent相关依赖:
- 下载与Schema Registry版本匹配的JAR包(建议用7.x系列):
kafka-avro-serializer-<version>.jarkafka-schema-registry-client-<version>.jaravro-<version>.jarjackson-databind-<version>.jar(若缺失)
- 启动扩展服务时指定依赖路径:
java -cp "<beam-runners-java-fn-executor-2.44.0.jar>:<confluent-jars-path>/*" org.apache.beam.sdk.expansion.service.ExpansionService 8097
- 在Python代码中指定扩展服务地址:
options = PipelineOptions([ "--expansion_service=localhost:8097" ])
2. 显式指定Avro Coder
在ReadFromKafka后添加Coder配置,确保Beam能正确处理Avro数据:
from apache_beam.coders.avro_coder import AvroCoder # 替换为你的实际Avro Schema,可从Schema Registry获取 schema = """{ "type": "record", "name": "YourRecord", "fields": [{"name": "field1", "type": "string"}, {"name": "field2", "type": "int"}] }""" with beam.Pipeline(options=options) as pipeline: ( pipeline | 'Read' >> ReadFromKafka(...) | 'Set Avro Coder' >> beam.Map(lambda x: x).with_output_types(AvroCoder(schema)) | 'Print' >> beam.Map(print) )
3. 优化Kafka配置
- 移除冗余配置:
sasl.username和sasl.password已在sasl.jaas.config中定义,无需重复设置。 - 调整Schema Registry地址格式,单独配置认证信息:
conf['schema.registry.url'] = f"https://{host}:29650" conf['basic.auth.user.info'] = f"{username}:{password}"
4. M1平台适配
- 用Rosetta 2运行Java扩展服务,规避ARM架构兼容性问题:
arch -x86_64 java -cp "<jar-path>" org.apache.beam.sdk.expansion.service.ExpansionService 8097
- 创建x86_64架构的conda环境(通过
arch -x86_64 conda create命令),减少跨架构运行问题。
内容的提问来源于stack exchange,提问作者Maria Zorkaltseva
相关产品推荐
相关产品推荐

