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

使用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,具体原因包括:

  1. Java依赖缺失:Apache Beam的Kafka外部Transform依赖Confluent的Avro序列化库,但默认的Beam Java扩展服务未包含这些依赖,导致无法加载KafkaAvroDeserializer类。
  2. Coder未显式指定:使用Java端的Avro反序列化器时,Beam无法自动推断对应的Coder,导致数据格式解析失败。
  3. M1平台兼容性:ARM架构下部分Java依赖存在适配问题,可能引发类加载异常。

解决建议

1. 补充Java依赖并启动自定义扩展服务

Beam Python SDK调用Kafka外部Transform依赖Java扩展服务,需手动添加Confluent相关依赖:

  • 下载与Schema Registry版本匹配的JAR包(建议用7.x系列):
    • kafka-avro-serializer-<version>.jar
    • kafka-schema-registry-client-<version>.jar
    • avro-<version>.jar
    • jackson-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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 01:42:22