【发布时间】: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