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

AWS EMR上PyFlink作业执行失败:NoClassDefFoundError问题排查

问题描述

在AWS EMR v7.3.0集群上,使用Python 3.9和PyFlink运行Flink作业,从AWS Kinesis流读取数据并打印到控制台。通过SSH连接主节点执行作业时失败,核心错误为java.lang.NoClassDefFoundError: Could not initialize class org.apache.flink.kinesis.shaded.com.amazonaws.ClientConfiguration。

作业代码(test.py)

import argparse
import json
from pyflink.datastream.connectors.kinesis import FlinkKinesisConsumer
from pyflink.common import SimpleStringSchema, Types
from pyflink.datastream import StreamExecutionEnvironment, RuntimeExecutionMode
from pyflink.datastream.functions import MapFunction


class ProcessImage(MapFunction):
    def open(self, runtime_context):
        print("Constructor launched ...")

    def map(self, value):
        try:
            if not value:
                print("Empty Kinesis record received.")
                return None

            # Parse the JSON-formatted record
            record = json.loads(value)
            print(f"Processed record: {record}")

            return value
        except Exception as e:
            print(f"Error processing record: {e}")
            raise


if __name__ == "__main__":
    parser = argparse.ArgumentParser(description="StreamingApp")

    parser.add_argument(
        "--stream_name",
        type=str,
        help="The name of Kinesis Stream to connect Flink with.",
        required=True,
    )
    parser.add_argument(
        "--region",
        type=str,
        help="The region name of the streaming application.",
        choices=["us-east-1", "us-west-2"],
        required=True,
    )

    args = parser.parse_args()
    stream_name: str = args.stream_name
    region: str = args.region

    # Set up the Flink environment
    env = StreamExecutionEnvironment.get_execution_environment()
    env.set_runtime_mode(RuntimeExecutionMode.STREAMING)

    env.add_jars(
        "file:////home/hadoop/flink-connector-kinesis-4.3.0-1.18.jar",
        "file:////home/hadoop/joda-time-2.12.5.jar",
    )

    # Kinesis Consumer properties
    kinesis_consumer_config = {
        "aws.region": region,
        "stream.initial.position": "LATEST",
        "aws.credentials.provider": "AUTO",
    }

    # Set up the Kinesis consumer to read from the Kinesis stream
    kinesis_source = FlinkKinesisConsumer(
        stream_name,
        SimpleStringSchema(),
        kinesis_consumer_config,
    )

    # Define the stream pipeline
    stream = env.add_source(kinesis_source)

    # Process record
    processed_stream = stream.map(ProcessImage(), output_type=Types.STRING())
    processed_stream.print()

    # Execute the Flink job
    env.execute()

执行命令及错误栈

sh-5.2$ python3 /home/hadoop/test.py --stream_name TestStream --region us-west-2
Setting HADOOP_CONF_DIR=/etc/hadoop/conf because no HADOOP_CONF_DIR orHADOOP_CLASSPATH was set.
Setting HBASE_CONF_DIR=/etc/hbase/conf because no HBASE_CONF_DIR was set.
Traceback (most recent call last):
  File "/home/hadoop/test.py", line 67, in <module>
    main()
  File "/home/hadoop/test.py", line 63, in main
    env.execute("Flink Kinesis Processing Job")
  File "/home/ssm-user/.local/lib/python3.9/site-packages/pyflink/datastream/stream_execution_environment.py", line 824, in execute
    return JobExecutionResult(self._j_stream_execution_environment.execute(j_stream_graph))
  File "/home/ssm-user/.local/lib/python3.9/site-packages/py4j/java_gateway.py", line 1322, in __call__
    return_value = get_return_value(
  File "/home/ssm-user/.local/lib/python3.9/site-packages/pyflink/util/exceptions.py", line 146, in deco
    return f(*a, **kw)
  File "/home/ssm-user/.local/lib/python3.9/site-packages/py4j/protocol.py", line 326, in get_return_value
    raise Py4JJavaError(
py4j.protocol.Py4JJavaError: An error occurred while calling o0.execute.
: org.apache.flink.runtime.client.JobExecutionException: Job execution failed.
...
Caused by: org.apache.flink.runtime.JobException: Recovery is suppressed by NoRestartBackoffTimeStrategy
....
Caused by: java.lang.NoClassDefFoundError: Could not initialize class org.apache.flink.kinesis.shaded.com.amazonaws.ClientConfiguration
        at org.apache.flink.kinesis.shaded.com.amazonaws.ClientConfigurationFactory.getDefaultConfig(ClientConfigurationFactory.java:46)
        at org.apache.flink.kinesis.shaded.com.amazonaws.ClientConfigurationFactory.getConfig(ClientConfigurationFactory.java:36)
        at org.apache.flink.streaming.connectors.kinesis.proxy.KinesisProxy.createKinesisClient(KinesisProxy.java:268)
        at org.apache.flink.streaming.connectors.kinesis.proxy.KinesisProxy.<init>(KinesisProxy.java:152)
        at org.apache.flink.streaming.connectors.kinesis.proxy.KinesisProxy.create(KinesisProxy.java:280)
        at org.apache.flink.streaming.connectors.kinesis.internals.KinesisDataFetcher.<init>(KinesisDataFetcher.java:412)
        at org.apache.flink.streaming.connectors.kinesis.internals.KinesisDataFetcher.<init>(KinesisDataFetcher.java:366)
        at org.apache.flink.streaming.connectors.kinesis.FlinkKinesisConsumer.createFetcher(FlinkKinesisConsumer.java:541)
        at org.apache.flink.streaming.connectors.kinesis.FlinkKinesisConsumer.run(FlinkKinesisConsumer.java:308)
        at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:113)
        at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:71)
        at org.apache.flink.streaming.runtime.tasks.SourceStreamTask$LegacySourceFunctionThread.run(SourceStreamTask.java:338)

核心错误

Caused by: java.lang.NoClassDefFoundError: Could not initialize class org.apache.flink.kinesis.shaded.com.amazonaws.ClientConfiguration
        at org.apache.flink.kinesis.shaded.com.amazonaws.ClientConfigurationFactory.getDefaultConfig(ClientConfigurationFactory.java:46)

已尝试在Python脚本中通过env.add_jars()添加flink-connector-kinesis-4.3.0-1.18.jar和joda-time-2.12.5.jar,且IAM权限验证正确,作业仍失败。咨询以下问题:

  1. 是否有人遇到过此类问题?是否可能是AWS SDK与Flink Kinesis连接器版本不兼容导致?
  2. 在EMR环境中使用PyFlink的Kinesis连接器,是否需要额外依赖或配置?
  3. 该问题是否与EMR上Flink的默认类加载器配置有关?

解决方案与问题解答

1. 版本兼容性问题

这个错误确实和版本不兼容直接相关。EMR 7.3.0默认搭载的Flink版本是1.18.0,你使用的flink-connector-kinesis-4.3.0-1.18.jar虽然版本匹配,但EMR环境自带的AWS SDK依赖可能和连接器内置的shaded AWS SDK存在冲突。NoClassDefFoundError通常是类初始化失败导致,大概率是依赖版本冲突引发的类加载异常。

2. EMR环境的额外依赖与配置

在EMR上运行PyFlink+Kinesis作业,不能仅通过env.add_jars()添加连接器JAR,需注意以下几点:

  • 使用EMR预装的连接器:EMR已为Flink预装适配的Kinesis连接器,路径通常为/usr/lib/flink/lib/flink-connector-kinesis-*.jar,直接使用这个预装JAR,不要自行下载第三方版本,避免版本冲突。
  • 提交作业时指定参数:提交作业时需确保连接器的shaded类优先加载,可添加参数:
    python3 test.py --stream_name TestStream --region us-west-2 \
    --jarfile /usr/lib/flink/lib/flink-connector-kinesis-4.3.0-1.18.jar \
    -Dclassloader.resolve-order=parent-first
    
  • 验证依赖完整性:需确保flink-shaded-aws-sdk-*.jar存在于Flink的lib目录中,该JAR包含连接器依赖的shaded AWS SDK类,EMR通常会预装,若缺失需手动添加。

3. 类加载器配置的影响

是的,这个问题和EMR上Flink的类加载器配置有关。EMR默认的类加载器顺序可能导致系统自带的AWS SDK类优先于连接器的shaded类加载,从而引发初始化失败。通过设置-Dclassloader.resolve-order=parent-first可让连接器的shaded类优先加载,避免冲突。

也可以修改Flink全局配置文件/etc/flink/conf/flink-conf.yaml,添加:

classloader.resolve-order: parent-first

这样所有作业都会使用该类加载顺序,无需每次提交作业时指定参数。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 22:14:53