使用MongoDBtoBigQuery模板时Javascript UDF过滤触发NullPointerException
使用MongoDBtoBigquery批处理模板从MongoDB迁移数据到BigQuery时,只要用UDF过滤记录就会触发NullPointerException,不用UDF时模板运行正常。
源数据里虽然有可选字段,但用来过滤的status字段是必存在的。猜测模板要求UDF必须返回所有行,不能跳过记录,但这样就搞不懂这类模板里UDF的用途了,难道只能用来增删列?
想知道错误原因,以及能不能通过Dataflow模板结合UDF实现过滤操作。
所用UDF代码
/** * @param {string} inJson * @return {string} outJson */ function process(inJson) { var obj = JSON.parse(inJson); // 仅输出status为"deleted"的对象 if (obj.hasOwnProperty('status') && obj['status'] === "deleted") { return JSON.stringify(obj); } }
错误日志
Error message from worker: org.apache.beam.sdk.util.UserCodeException: java.lang.NullPointerException
org.apache.beam.sdk.util.UserCodeException.wrap(UserCodeException.java:39)
com.google.cloud.teleport.v2.mongodb.templates.MongoDbToBigQuery$1$DoFnInvoker.invokeProcessElement(Unknown Source)
org.apache.beam.fn.harness.FnApiDoFnRunner.processElementForParDo(FnApiDoFnRunner.java:803)
org.apache.beam.fn.harness.data.PCollectionConsumerRegistry$MetricTrackingFnDataReceiver.accept(PCollectionConsumerRegistry.java:348)
org.apache.beam.fn.harness.data.PCollectionConsumerRegistry$MetricTrackingFnDataReceiver.accept(PCollectionConsumerRegistry.java:275)
org.apache.beam.fn.harness.FnApiDoFnRunner.outputTo(FnApiDoFnRunner.java:1792)
org.apache.beam.fn.harness.FnApiDoFnRunner.access$3000(FnApiDoFnRunner.java:143)
org.apache.beam.fn.harness.FnApiDoFnRunner$NonWindowObservingProcessBundleContext.output(FnApiDoFnRunner.java:2514)
com.google.cloud.teleport.v2.transforms.JavascriptDocumentTransformer$TransformDocumentViaJavascript$1.processElement(JavascriptDocumentTransformer.java:237)
com.google.cloud.teleport.v2.transforms.JavascriptDocumentTransformer$TransformDocumentViaJavascript$1$DoFnInvoker.invokeProcessElement(Unknown Source)
org.apache.beam.fn.harness.FnApiDoFnRunner.processElementForParDo(FnApiDoFnRunner.java:803)
org.apache.beam.fn.harness.data.PCollectionConsumerRegistry$MetricTrackingFnDataReceiver.accept(PCollectionConsumerRegistry.java:348)
org.apache.beam.fn.harness.data.PCollectionConsumerRegistry$MetricTrackingFnDataReceiver.accept(PCollectionConsumerRegistry.java:275)
org.apache.beam.fn.harness.FnApiDoFnRunner.outputTo(FnApiDoFnRunner.java:1792)
org.apache.beam.fn.harness.FnApiDoFnRunner.access$3000(FnApiDoFnRunner.java:143)
org.apache.beam.fn.harness.FnApiDoFnRunner$WindowObservingProcessBundleContext.outputWithTimestamp(FnApiDoFnRunner.java:2218)
org.apache.beam.sdk.io.Read$BoundedSourceAsSDFWrapperFn.processElement(Read.java:321)
org.apache.beam.sdk.io.Read$BoundedSourceAsSDFWrapperFn$DoFnInvoker.invokeProcessElement(Unknown Source)
org.apache.beam.fn.harness.FnApiDoFnRunner.processElementForWindowObservingSizedElementAndRestriction(FnApiDoFnRunner.java:1100)
org.apache.beam.fn.harness.FnApiDoFnRunner.access$1500(FnApiDoFnRunner.java:143)
org.apache.beam.fn.harness.FnApiDoFnRunner$4.accept(FnApiDoFnRunner.java:659)
org.apache.beam.fn.harness.FnApiDoFnRunner$4.accept(FnApiDoFnRunner.java:654)
org.apache.beam.fn.harness.data.PCollectionConsumerRegistry$MetricTrackingFnDataReceiver.accept(PCollectionConsumerRegistry.java:348)
org.apache.beam.fn.harness.data.PCollectionConsumerRegistry$MetricTrackingFnDataReceiver.accept(PCollectionConsumerRegistry.java:275)
org.apache.beam.fn.harness.BeamFnDataReadRunner.forwardElementToConsumer(BeamFnDataReadRunner.java:213)
org.apache.beam.sdk.fn.data.BeamFnDataInboundObserver.multiplexElements(BeamFnDataInboundObserver.java:158)
org.apache.beam.fn.harness.control.ProcessBundleHandler.processBundle(ProcessBundleHandler.java:537)
org.apache.beam.fn.harness.control.BeamFnControlClient.delegateOnInstructionRequestType(BeamFnControlClient.java:150)
org.apache.beam.fn.harness.control.BeamFnControlClient$InboundObserver.lambda$onNext$0(BeamFnControlClient.java:115)
java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
org.apache.beam.sdk.util.UnboundedScheduledExecutorService$ScheduledFutureTask.run(UnboundedScheduledExecutorService.java:163)
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.NullPointerException
com.google.cloud.teleport.v2.mongodb.templates.MongoDbUtils.getTableSchema(MongoDbUtils.java:119)
com.google.cloud.teleport.v2.mongodb.templates.MongoDbToBigQuery$1.process(MongoDbToBigQuery.java:146)
错误原因
从错误栈可以看出,空指针异常触发自MongoDbUtils.getTableSchema方法。问题出在你的UDF上:当记录不满足过滤条件时,函数没有返回值(默认返回undefined),而该Dataflow模板的UDF执行逻辑期望每个输入都对应一个非空的JSON字符串输出。当UDF返回undefined时,后续处理BigQuery表结构的代码因为拿到空值,直接抛出了NullPointerException。
这个模板的UDF设计初衷确实是仅用于转换数据结构(增删字段、修改字段值),而非过滤记录——它不会自动忽略UDF返回的空值,而是会把空值传入下游逻辑,导致报错。
实现过滤的可行方案
方案1:修改UDF,配合后续过滤步骤
先让UDF给不需要保留的记录标记一个特殊标识,然后在Dataflow管道中添加过滤步骤移除这些记录:
/** * @param {string} inJson * @return {string} outJson */ function process(inJson) { var obj = JSON.parse(inJson); // 给不符合条件的记录添加标记字段 if (!(obj.hasOwnProperty('status') && obj['status'] === "deleted")) { obj._skip = true; } return JSON.stringify(obj); }
之后如果是自定义模板,在后续逻辑中添加ParDo过滤掉_skip: true的记录;如果使用官方预制模板无法修改代码,可选择方案2。
方案2:在MongoDB端先过滤
直接在MongoDB的读取阶段添加查询条件,只读取status: "deleted"的记录,这样不需要在Dataflow中做过滤,避免UDF引发的问题。
方案3:自定义Dataflow管道
如果预制模板无法满足需求,建议基于官方模板代码修改:
- 在UDF执行完成后,添加过滤步骤,检查UDF输出是否为非空值,只保留有效输出;
- 或者将过滤逻辑和数据转换逻辑分开,UDF只负责结构转换,过滤交给专门的ParDo步骤处理。
内容的提问来源于stack exchange,提问作者Dchemist_Rae

