【问题标题】:Slowness / Lag in beam streaming pipeline in group by key stage关键阶段分组流束流水线的慢度/滞后
【发布时间】:2020-03-31 20:13:40
【问题描述】:

上下文

大家好,我一直在使用Apache Beam 管道生成柱状数据库以存储在GCS 中,我有一个来自Kafka 的数据流并且有一个1m 的窗口。

我想将那个 1m 窗口的所有数据转换成一个柱状 DB 文件(在我的例子中是 ORC,可以是 Parquet 或其他任何东西),我已经为这个转换编写了一个管道。

问题

我正在经历普遍的缓慢。我怀疑这可能是由于按键组转换,因为我只有按键。真的有必要这样做吗?如果没有,应该怎么做?我读到 combine 对此并不是很有用,因为我的管道并没有真正聚合数据,而是创建了一个合并文件。我真正需要的是每个窗口的可迭代对象列表,它将被转换为 ORC 文件。

管道表示

输入->窗口->按键分组(只有1个键)-> pardo(创建DB)-> IO(写入GCS)

我尝试过的

我尝试过使用分析器,水平/垂直缩放。使用分析器,我看到超过 50% 的时间通过按键操作分组。我确实相信问题出在热键上,但我无法找到应该做什么的解决方案。当我删除group by key 操作时,我的管道跟上Kafka 的滞后(即,在Kafka 结束时这似乎不是问题)。

代码片段

p.apply("ReadLines", KafkaIO.<Long, byte[]>read().withBootstrapServers("myserver.com:9092")
    .withTopic(options.getInputTopic())
    .withTimestampPolicyFactory(MyTimePolicy.myTimestampPolicyFactory())
    .withConsumerConfigUpdates(Map.of("group.id", "mygroup-id")).commitOffsetsInFinalize()
    .withKeyDeserializer(LongDeserializer.class)
    .withValueDeserializer(ByteArrayDeserializer.class).withoutMetadata())
    .apply("UncompressSnappy", ParDo.of(new UncompressSnappy()))
    .apply("DecodeProto", ParDo.of(new DecodePromProto()))
    .apply("MapTSSample", ParDo.of(new MapTSSample()))
    .apply(Window.<TSSample>into(FixedWindows.of(Duration.standardMinutes(1)))
        .withTimestampCombiner(TimestampCombiner.END_OF_WINDOW))
    .apply(WithKeys.<Integer, TSSample>of(1))
    .apply(GroupByKey.<Integer, TSSample>create())
    .apply("CreateTSORC", ParDo.of(new CreateTSORC()))
    .apply(new WriteOneFilePerWindow(options.getOutput(), 1));

墙上时间配置文件

https://gist.github.com/anandsinghkunwar/4cc26f7e3da7473af66ce9a142a74c35

【问题讨论】:

  • Beam 默认使用动态延迟配置,试图考虑它认为您的系统具有的延迟时间。我注意到使用 GCP Pub/Sub 作为源,如果我不进行窗口化,我的数据会立即通过我的整个 Beam 管道,但如果我使用窗口化,则在窗口关闭后会有明显的延迟。我相信它在看到窗口时使用了明显的延迟配置。您是否尝试过调整允许的延迟设置,或者手动将其设置为非常低的值?
  • 您使用的是 Dataflow 还是 DirectRunner?你能提供你的代码的sn-p吗?可能和youe窗口的一些配置有关
  • @rmesteves 添加了代码 sn-p。我正在使用数据流。
  • @MattWelke 我相信默认允许的迟到是 0,我应该明确将其设置为 0 吗?
  • @MattWelke 我已尝试将其显式设置为 0,但不会导致任何更改。

标签: google-cloud-dataflow apache-beam


【解决方案1】:

问题确实似乎是hot keys issue,我不得不更改我的管道来为 ORC 文件创建自定义 IO,并将分片数量增加到50。我完全删除了GroupByKey。由于 beam 还没有自动确定 FileIO.write() 的分片数量,因此您必须手动选择适合您工作量的数字。

此外,在 Google Dataflow 中启用流引擎 API 可以进一步加快提取速度。

【讨论】:

    猜你喜欢
    • 2019-02-19
    • 1970-01-01
    • 2020-11-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-07
    相关资源
    最近更新 更多