【问题标题】:Can Dataflow sideInput be updated per window by reading a gcs bucket?可以通过读取 gcs 存储桶来更新每个窗口的 Dataflow sideInput 吗?
【发布时间】:2017-05-06 08:39:11
【问题描述】:

我目前正在创建一个 PCollectionView,方法是从 gcs 存储桶中读取过滤信息,并将其作为侧输入传递到管道的不同阶段,以过滤输出。如果 gcs 存储桶中的文件发生更改,我希望当前正在运行的管道使用这个新的过滤器信息。如果我的过滤器发生变化,有没有办法在每个新的数据窗口上更新这个 PCollectionView?我以为我可以在 startBundle 中做到这一点,但我不知道如何或是否可能。可以的话可以举个例子吗?

PCollectionView<Map<String, TagObject>> 
    tagMapView =
        pipeline.apply(TextIO.Read.named("TagListTextRead")
                                  .from("gs://tag-list-bucket/tag-list.json"))
                .apply(ParDo.named("TagsToTagMap").of(new Tags.BuildTagListMapFn()))
                .apply("MakeTagMapView", View.asSingleton());
PCollection<String> 
    windowedData =
        pipeline.apply(PubsubIO.Read.topic("myTopic"))
                .apply(Window.<String>into(
                              SlidingWindows.of(Duration.standardMinutes(15))
                                            .every(Duration.standardSeconds(31))));
PCollection<MY_DATA> 
    lineData = windowedData
        .apply(ParDo.named("ExtractJsonObject")
            .withSideInputs(tagMapView)
            .of(new ExtractJsonObjectFn()));

【问题讨论】:

    标签: google-cloud-dataflow


    【解决方案1】:

    您可能想要“使用最多 1 分钟旧版本的过滤器作为辅助输入”之类的东西(因为从理论上讲,文件可以频繁、不可预测地更改并且独立于您的管道 - 所以没有办法真正将文件的更改与管道的行为完全同步)。

    这是我能够想出的(公认的,相当笨拙的)解决方案。它依赖于这样一个事实,即侧面输入也由窗口隐式键入。在这个解决方案中,我们将创建一个以 1 分钟固定窗口为窗口的侧面输入,其中每个窗口将包含标签映射的单个值,该值源自该窗口内的某个时刻的过滤器文件。

    PCollection<Long> ticks = p
      // Produce 1 "tick" per second
      .apply(CountingInput.unbounded().withRate(1, Duration.standardSeconds(1)))
      // Window the ticks into 1-minute windows
      .apply(Window.into(FixedWindows.of(Duration.standardMinutes(1))))
      // Use an arbitrary per-window combiner to reduce to 1 element per window
      .apply(Count.globally());
    
    // Produce a collection of tag maps, 1 per each 1-minute window
    PCollectionView<TagMap> tagMapView = ticks
      .apply(MapElements.via((Long ignored) -> {
        ... manually read the json file as a TagMap ...
      }))
      .apply(View.asSingleton());
    

    这种模式(将缓慢变化的外部数据作为辅助输入加入)反复出现,我在这里提出的解决方案远非完美,我希望我们在编程模型中能更好地支持这一点。我已经提交了BEAM JIRA issue 来跟踪这个。

    【讨论】:

    • 非常感谢您的解决方案,但我无法从“忽略”中获得正确的填写语法。我不确定 MapElements.via 想要什么。能否请您更准确地填写一下。
    • 这是一个 Java 8 lambda,但实际上我的语法错误。确保您使用的是 Java 8 - 否则,请使用常规 ParDo 和匿名类。
    • 上面的答案似乎有问题。 Count.globally() 需要执行 GlobalWindow,这意味着它不能特别用于流数据流。这里有没有遗漏的部分?
    • CountingInput 似乎不再存在。 (我使用的是 Beam SDK 2.6)
    • 现在叫做 GenerateSequence。
    猜你喜欢
    • 1970-01-01
    • 2019-11-01
    • 2015-04-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多