使用Flink执行Hudi Compaction时触发NullPointerException求助
问题描述
参考Hudi官网的Flink离线Compaction示例操作,未使用示例中的hudi-flink-bundle_2.11-0.9.0-SNAPSHOT.jar,改用Maven仓库获取的hudi-flink1.16-bundle-0.13.0.jar,执行命令如下:
$FLINK_HOME/bin/flink run \ -c org.apache.hudi.sink.compact.HoodieFlinkCompactor \ $FLINK_HOME/lib/hudi-flink1.16-bundle-0.13.0.jar \ --path 'file://...sample_db/people'
运行后触发NullPointerException,异常栈信息如下:
The program finished with the following exception: org.apache.flink.client.program.ProgramInvocationException: The main method caused an error: Value must not be null. at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:372) at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:222) at org.apache.flink.client.ClientUtils.executeProgram(ClientUtils.java:98) at org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:843) at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:240) at org.apache.flink.client.cli.CliFrontend.parseAndRun(CliFrontend.java:1087) at org.apache.flink.client.cli.CliFrontend.lambda$main$10(CliFrontend.java:1165) at java.base/java.security.AccessController.doPrivileged(Native Method) at java.base/javax.security.auth.Subject.doAs(Subject.java:423) at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1878) at org.apache.flink.runtime.security.contexts.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41) at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1165) Caused by: java.lang.NullPointerException: Value must not be null. at org.apache.flink.configuration.Configuration.setValueInternal(Configuration.java:775) at org.apache.flink.configuration.Configuration.setValueInternal(Configuration.java:787) at org.apache.flink.configuration.Configuration.setString(Configuration.java:200) at org.apache.hudi.util.CompactionUtil.setPreCombineField(CompactionUtil.java:133) at org.apache.hudi.sink.compact.HoodieFlinkCompactor$AsyncCompactionService.<init>(HoodieFlinkCompactor.java:177) at org.apache.hudi.sink.compact.HoodieFlinkCompactor.main(HoodieFlinkCompactor.java:75) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:566) at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:355)
请问是否有人遇到过该情况?可能的根因是什么?
可能的根因及解决方法
- 根因:从异常栈可以定位到,
CompactionUtil.setPreCombineField方法尝试往Flink配置中写入预合并字段值,但该值为null。这是因为Hudi 0.13.0版本的离线Compaction逻辑要求必须指定预合并字段,要么通过命令行参数传递,要么表元数据中已存储该字段信息,二者缺一不可。 - 解决方法:
- 在执行Compaction命令时,显式添加
--precombine-field参数,指定对应表的预合并字段(比如时间戳字段ts):$FLINK_HOME/bin/flink run \ -c org.apache.hudi.sink.compact.HoodieFlinkCompactor \ $FLINK_HOME/lib/hudi-flink1.16-bundle-0.13.0.jar \ --path 'file://...sample_db/people' \ --precombine-field 'ts' - 检查目标Hudi表的创建元数据,确认是否配置了
hoodie.datasource.write.precombine.field参数。如果表创建时未配置该参数,后续所有写入和Compaction操作都需要显式传递预合并字段参数。
- 在执行Compaction命令时,显式添加
内容的提问来源于stack exchange,提问作者Bing-hsu Gao
相关产品推荐
相关产品推荐

