【问题标题】:NACK not send back to Google Cloud Pub/Sub from Dataflow when there is an error in ParDo function当 ParDo 函数出现错误时,NACK 不会从 Dataflow 发送回 Google Cloud Pub/Sub
【发布时间】:2021-07-16 06:40:45
【问题描述】:

当 Dataflow 作业无法或​​不愿处理消息时,如何向 Pub/Sub 发送 NACK。

Pipeline pipeline = Pipeline.create(options);

    pipeline.apply("gcs2ZipExtractor-processor",
            PubsubIO.readMessagesWithAttributes()
                    .fromSubscription(pubSubSubscription))
           .apply(ParDo.of(new ProcessZipFileEventDoFn(appProps)));
    logger.info("Started ZipFile Extractor");
    pipeline.run().waitUntilFinish();

上面是我用来运行 ApacheBeam Dataflow 管道作业的代码 sn-p。如果 ProcessZipFileEventDoFn 发生任何故障,我想向 Pub/Sub 订阅发送 NACK 消息,以便将消息移动到 DeadletterTopic。目前,Dataflow Runner 没有发生 NACK。

【问题讨论】:

    标签: apache-beam google-cloud-pubsub dataflow


    【解决方案1】:

    目前,Apache Beam SDK 不支持 Pub/Sub 的原生死信队列功能。但是,您可以相当轻松地编写自己的代码。以下内容来自此blog post 改编为您的代码。诀窍是使用来自单个 ParDo 的多个输出。一个输出 PCollection 将具有不引发任何异常的“好”数据。如果有任何异常,另一个输出 PCollection 将包含所有“坏”数据。然后,您可以将死信 PCollection 中的所有元素写入接收器,在您的情况下为 Pub/Sub 主题。

    PCollection input =
        pipeline.apply("gcs2ZipExtractor-processor",
                       PubsubIO.readMessagesWithAttributes()
                           .fromSubscription(pubSubSubscription))
    
    // Put this try-catch logic in your ProcessZipFileEventDoFn, and don't forget
    // the "withOutputTags"!
    final TupleTag successTag ;
    final TupleTag deadLetterTag;
    PCollectionTuple outputTuple = input.apply(ParDo.of(new DoFn() {  
      @Override  
      void processElement(ProcessContext c) {    
      try {      
        c.output(process(c.element());    
      } catch (Exception e) {      
        c.sideOutput(deadLetterTag, c.element());  
      }
    }).withOutputTags(successTag, TupleTagList.of(deadLetterTag)));
    
    // Write the dead letter inputs to Pub/Sub for later analysis
    outputTuple.get(deadLetterTag).apply(PubSubIO.write(...));
    
    // Retrieve the successful elements...
    PCollection success = outputTuple.get(successTag);
    

    【讨论】:

      猜你喜欢
      • 2020-02-26
      • 2015-11-09
      • 2021-02-21
      • 2019-11-29
      • 1970-01-01
      • 2019-10-12
      • 2022-01-13
      • 2017-09-14
      • 2017-11-27
      相关资源
      最近更新 更多