【问题标题】:apache beam's StartBundle is filing on weird errorapache beam StartBundle 正在归档奇怪的错误
【发布时间】:2021-04-29 12:13:23
【问题描述】:

我收到了这个错误

Caused by: java.lang.IllegalArgumentException: 
com.orderly.rosters.transforms.RosterFileReader$RosterFileReaderFn, @StartBundle start(StartBundleContext), parameter of type StartBundleContext at index 0: StartBundleContext argument must have type DoFn<String, List<String>>.ProcessContext

关于以下代码

public abstract class OrderlyDoFn<INPUT, OUTPUT> extends DoFn<INPUT, OUTPUT> {
    protected Logger log = LoggerFactory.getLogger(getClass());
    private transient String projectId;
    private transient Map headers;

    @DoFn.StartBundle
    public void start(DoFn.StartBundleContext ctx) {
        OrderlyPipelineOptions options = (OrderlyPipelineOptions) ctx.getPipelineOptions();
        headers = PlatformMagic.unmarshal(options.getPlatformMagic().get(), Map.class);
    }

    @DoFn.ProcessElement
    public void processElement(@DoFn.Element INPUT elem, DoFn.OutputReceiver<OUTPUT> receiver) {
        try {
            RouterRequest routerReq = new RouterRequest();
            routerReq.requestState = headers;
            RequestContext ctx = new RequestContext(null, null, null, routerReq, null);
            Current.setContext(ctx);
            routerReq.requestState.put(DataflowClientFactory.PROJECT_KEY, projectId);

            for (OrderlyHeaders header : OrderlyHeaders.values()) {
                if (header.isLogged()) {
                    String value = (String) routerReq.requestState.get(header.getHeaderName());
                    MDC.put(header.getLoggerKey(), value);
                }
            }

            processElementImpl(elem, receiver);
        } catch (Throwable e) {
            log.info("Exception processing OrderlyDoFn", e);
            throw SneakyThrow.sneak(e);
        } finally {
            MDC.clear();
            Current.setContext(null); //clear context
        }
    }

    protected abstract void processElementImpl(INPUT elem, OutputReceiver<OUTPUT> receiver);

}

【问题讨论】:

    标签: google-cloud-dataflow apache-beam


    【解决方案1】:

    呵呵,解决了

    @DoFn.StartBundle
    public void start(DoFn<INPUT, OUTPUT>.StartBundleContext ctx) {
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-08-03
      • 2012-08-28
      • 1970-01-01
      • 2013-01-18
      • 1970-01-01
      • 2021-03-09
      • 1970-01-01
      • 2021-11-20
      相关资源
      最近更新 更多