如何在Apache Flink集群作业中正确传递--add-opens JVM参数
解决Flink作业在Java 17下提交时的模块访问错误
问题背景
使用Java 17和Apache Flink 1.18.1开发的Flink作业,提交到集群时出现以下错误:
java.lang.RuntimeException: Exception occurred while setting the current key context. at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.setCurrentKey(StreamOperatorStateHandler.java:373) at org.apache.flink.streaming.api.operators.AbstractStreamOperator.setCurrentKey(AbstractStreamOperator.java:508) at org.apache.flink.streaming.api.operators.AbstractStreamOperator.setKeyContextElement(AbstractStreamOperator.java:503) at org.apache.flink.streaming.api.operators.AbstractStreamOperator.setKeyContextElement1(AbstractStreamOperator.java:478) at org.apache.flink.streaming.api.operators.OneInputStreamOperator.setKeyContextElement(OneInputStreamOperator.java:36) at org.apache.flink.streaming.runtime.io.RecordProcessorUtils.lambda$getRecordProcessor$0(RecordProcessorUtils.java:59) at org.apache.flink.streaming.runtime.tasks.OneInputStreamTask$StreamTaskNetworkOutput.emitRecord(OneInputStreamTask.java:237) at org.apache.flink.streaming.runtime.io.AbstractStreamTaskNetworkInput.processElement(AbstractStreamTaskNetworkInput.java:146) at org.apache.flink.streaming.runtime.io.AbstractStreamTaskNetworkInput.emitNext(AbstractStreamTaskNetworkInput.java:110) at org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:65) at org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:562) at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:231) at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:858) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:807) at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:953) at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:932) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:746) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:562) at java.base/java.lang.Thread.run(Unknown Source) Caused by: java.lang.RuntimeException: java.lang.reflect.InaccessibleObjectException: Unable to make field private final java.lang.Object[] java.util.Arrays$ArrayList.a accessible: module java.base does not "opens java.util" to unnamed module @24b6c462 at com.twitter.chill.java.ArraysAsListSerializer.<init>(ArraysAsListSerializer.java:69) at org.apache.flink.api.java.typeutils.runtime.kryo.FlinkChillPackageRegistrar.registerSerializers(FlinkChillPackageRegistrar.java:67) at org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializer.getKryoInstance(KryoSerializer.java:513) at org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializer.checkKryoInitialized(KryoSerializer.java:522) at org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializer.serialize(KryoSerializer.java:348) at org.apache.flink.runtime.state.SerializedCompositeKeyBuilder.serializeKeyGroupAndKey(SerializedCompositeKeyBuilder.java:192) at org.apache.flink.runtime.state.SerializedCompositeKeyBuilder.setKeyAndKeyGroup(SerializedCompositeKeyBuilder.java:95) at org.apache.flink.contrib.streaming.state.RocksDBKeyedStateBackend.setCurrentKey(RocksDBKeyedStateBackend.java:431) at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.setCurrentKey(StreamOperatorStateHandler.java:371) ... 18 more Caused by: java.lang.reflect.InaccessibleObjectException: Unable to make field private final java.lang.Object[] java.util.Arrays$ArrayList.a accessible: module java.base does not "opens java.util" to unnamed module @24b6c462 at java.base/java.lang.reflect.AccessibleObject.checkCanSetAccessible(Unknown Source) at java.base/java.lang.reflect.AccessibleObject.checkCanSetAccessible(Unknown Source) at java.base/java.lang.reflect.Field.checkCanSetAccessible(Unknown Source) at java.base/java.lang.reflect.Field.setAccessible(Unknown Source) at com.twitter.chill.java.ArraysAsListSerializer.<init>(ArraysAsListSerializer.java:67) ... 26 more
已确认需要以下JVM参数解决该问题:
--add-opens java.base/java.util=ALL-UNNAMED --add-opens java.base/java.lang=ALL-UNNAMED --add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED
尝试过以下方式但未生效:
- 通过
flink run传递客户端参数:
flink run -Djava.arg=--add-opens=java.base/java.util=ALL-UNNAMED -Djava.arg=--add-opens=java.base/java.lang=ALL-UNNAMED -Djava.arg=--add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED
- 在maven-shade-plugin的Manifest中添加配置:
<transformers> <transformer implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"/> <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"><mainClass>data.PubSubToGcsCsv</mainClass> <manifestEntries> <add-opens>java.base/java.util=ALL-UNNAMED</add-opens> <add-opens>java.base/java.lang=ALL-UNNAMED</add-opens> <add-opens>java.base/java.util.concurrent.atomic=ALL-UNNAMED</add-opens> </manifestEntries> </transformer> </transformers>
正确解决方案
方法1:提交作业时指定TaskManager的JVM参数
错误发生在TaskManager进程中,需将参数传递给TaskManager而非客户端:
flink run -Dtaskmanager.jvm.args="--add-opens java.base/java.util=ALL-UNNAMED --add-opens java.base/java.lang=ALL-UNNAMED --add-opens java.base/java.util.concurrent.atomic=ALL-UNNAMED" your-job.jar
注意:参数值需用引号包裹,避免空格导致参数拆分错误。
方法2:修改Flink集群配置(全局生效)
若所有作业都需要该权限,可修改集群的flink-conf.yaml文件:
taskmanager.jvm.args: "--add-opens java.base/java.util=ALL-UNNAMED --add-opens java.base/java.lang=ALL-UNNAMED --add-opens java.base/java.util.concurrent.atomic=ALL-UNNAMED"
修改后需重启TaskManager节点。
方法3:同时配置客户端与TaskManager参数
若客户端进程(如本地运行、客户端序列化阶段)也需要该权限,可同时添加客户端参数:
flink run -Dclient.jvm.args="--add-opens java.base/java.util=ALL-UNNAMED --add-opens java.base/java.lang=ALL-UNNAMED --add-opens java.base/java.util.concurrent.atomic=ALL-UNNAMED" -Dtaskmanager.jvm.args="--add-opens java.base/java.util=ALL-UNNAMED --add-opens java.base/java.lang=ALL-UNNAMED --add-opens java.base/java.util.concurrent.atomic=ALL-UNNAMED" your-job.jar
无效原因说明
-Djava.arg仅作用于Flink客户端进程,而错误发生在TaskManager中,因此无法解决问题。- Maven Shade插件的Manifest配置仅对直接运行JAR的进程生效,Flink集群的TaskManager通过自身启动脚本加载作业JAR,不会读取作业JAR的Manifest参数。
内容的提问来源于stack exchange,提问作者SRJ
相关产品推荐
相关产品推荐

