Streamsets Groovy脚本数据源报错求助:找不到state属性
StreamSets Groovy脚本数据源错误排查
问题场景
作为StreamSets和Groovy新手,将输入源从「Dev Raw Data Source」切换为「Groovy Scripting」后,执行以下脚本出现错误:
原Groovy脚本
import com.streamsets.pipeline.stage.origin.scripting.* record = sdc.createRecord("bulkUpdateRequestTemplate") bulkUpdateRequestTemplate = "{ \"payload\": { \"objects\": { \"filter\": \"(equals(type,'configuration/entityTypes/Contract') and equals(attributes.ContractStatus, 'ACTIVE') and lt (attributes.ExpirationDate,'currentTimestampEpoch'))\", \"options\": \"searchByOv,ovOnly\" }, \"actions\": [ { \"operation\": \"UpdateAttribute\", \"operationParameters\": { \"attributeURI\": \"configuration/entityTypes/Contract/attributes/ContractStatus\", \"attributeValue\": \"INACTIVE\" } } ] } }" record.value = bulkUpdateRequestTemplate sdc.state['bulkUpdateRequestTemplate'] = record
错误信息
com.streamsets.pipeline.api.StageException: SCRIPTING_10 - Script error in user script: javax.script.ScriptException: groovy.lang.MissingPropertyException: No such property: state for class: com.streamsets.pipeline.stage.origin.scripting.ScriptingOriginBindings at com.streamsets.pipeline.stage.origin.scripting.AbstractScriptingSource.produce(AbstractScriptingSource.java:124) at com.streamsets.pipeline.api.base.configurablestage.DPushSource.produce(DPushSource.java:44) at com.streamsets.datacollector.runner.StageRuntime.lambda$execute$1(StageRuntime.java:258) at com.streamsets.datacollector.runner.StageRuntime.execute(StageRuntime.java:232) at com.streamsets.datacollector.runner.StageRuntime.execute(StageRuntime.java:267) at com.streamsets.datacollector.runner.SourcePipe.process(SourcePipe.java:67) at com.streamsets.datacollector.runner.preview.PreviewPipelineRunner.runPushSource(PreviewPipelineRunner.java:233) at com.streamsets.datacollector.runner.preview.PreviewPipelineRunner.run(PreviewPipelineRunner.java:218) at com.streamsets.datacollector.runner.Pipeline.run(Pipeline.java:535) at com.streamsets.datacollector.runner.preview.PreviewPipeline.run(PreviewPipeline.java:39) at com.streamsets.datacollector.execution.preview.sync.SyncPreviewer.start(SyncPreviewer.java:227) at com.streamsets.datacollector.execution.preview.async.AsyncPreviewer.lambda$start$1(AsyncPreviewer.java:93) at com.streamsets.pipeline.lib.executor.SafeScheduledExecutorService$SafeCallable.lambda$call$0(SafeScheduledExecutorService.java:214) at com.streamsets.datacollector.security.GroupsInScope.execute(GroupsInScope.java:44) at com.streamsets.datacollector.security.GroupsInScope.execute(GroupsInScope.java:25) at com.streamsets.pipeline.lib.executor.SafeScheduledExecutorService$SafeCallable.call(SafeScheduledExecutorService.java:210) at com.streamsets.pipeline.lib.executor.SafeScheduledExecutorService$SafeCallable.lambda$call$0(SafeScheduledExecutorService.java:214) at com.streamsets.datacollector.security.GroupsInScope.execute(GroupsInScope.java:44) at com.streamsets.datacollector.security.GroupsInScope.execute(GroupsInScope.java:25) at com.streamsets.pipeline.lib.executor.SafeScheduledExecutorService$SafeCallable.call(SafeScheduledExecutorService.java:210) at java.util.concurrent.FutureTask.run(FutureTask.java:266) at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$201(ScheduledThreadPoolExecutor.java:180) at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:293) at com.streamsets.datacollector.metrics.MetricSafeScheduledExecutorService$MetricsTask.run(MetricSafeScheduledExecutorService.java:88) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:748)
错误原因
Scripting Origin(Groovy脚本数据源)的绑定对象中不存在sdc.state属性——sdc.state仅用于Processor(处理器)或Destination(目标端)阶段的状态持久化,源端没有这个功能。脚本试图将记录存入sdc.state,属于错误用法。
源端Groovy脚本的核心作用是生成并输出记录,正确方式是使用sdc.output.write()方法将记录发送到下游。
修正后的脚本
import com.streamsets.pipeline.stage.origin.scripting.* // 用Groovy多行字符串简化JSON写法,避免转义符 def bulkUpdateRequestTemplate = ''' { "payload": { "objects": { "filter": "(equals(type,'configuration/entityTypes/Contract') and equals(attributes.ContractStatus, 'ACTIVE') and lt (attributes.ExpirationDate,'currentTimestampEpoch'))", "options": "searchByOv,ovOnly" }, "actions": [ { "operation": "UpdateAttribute", "operationParameters": { "attributeURI": "configuration/entityTypes/Contract/attributes/ContractStatus", "attributeValue": "INACTIVE" } } ] } } ''' def record = sdc.createRecord("bulkUpdateRequestTemplate") record.value = bulkUpdateRequestTemplate // 输出记录到下游 sdc.output.write(record)
额外优化说明
- 用Groovy的**多行字符串(
'''包裹)**替代转义字符串,提升脚本可读性和维护性 - 移除无效的
sdc.state操作,改用sdc.output.write()完成记录输出 - 给变量添加
def关键字,符合Groovy最佳编程实践
内容的提问来源于stack exchange,提问作者Mike
相关产品推荐
相关产品推荐

