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权限验证正确,作业仍失败。咨询以下问题:
- 是否有人遇到过此类问题?是否可能是AWS SDK与Flink Kinesis连接器版本不兼容导致?
- 在EMR环境中使用PyFlink的Kinesis连接器,是否需要额外依赖或配置?
- 该问题是否与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

