【发布时间】:2020-08-12 06:36:58
【问题描述】:
我想在 GCP 中创建一个流式 Apache Beam 管道,它从 Google Pub/Sub 读取数据并将其推送到 GCS。我有一点可以从 Pub/Sub 读取数据。 我当前的代码看起来像这样(从 GCP Apache 梁模板之一中提取)
pipeline.apply("Read PubSub Events",
PubsubIO.readMessagesWithAttributes().fromTopic(options.getInputTopic()))
.apply("Map to Archive", ParDo.of(new PubsubMessageToArchiveDoFn()))
.apply(
options.getWindowDuration() + " Window",
Window.into(FixedWindows.of(DurationUtils.parseDuration(options.getWindowDuration()))))
.apply(
"Write File(s)",
AvroIO.write(AdEvent.class)
.to(
new WindowedFilenamePolicy(
options.getOutputDirectory(),
options.getOutputFilenamePrefix(),
options.getOutputShardTemplate(),
options.getOutputFilenameSuffix()))
.withTempDirectory(NestedValueProvider.of(
options.getAvroTempDirectory(),
(SerializableFunction<String, ResourceId>) input ->
FileBasedSink.convertToFileResourceIfPossible(input)))
.withWindowedWrites()
.withNumShards(options.getNumShards()));
它可以生成如下所示的文件
windowed-file2020-04-28T09:00:00.000Z-2020-04-28T09:02:00.000Z-pane-0-last-00-of-01.avro
我想将 GCS 中的数据存储在动态创建的目录中。在以下目录中,2020-04-28/01、2020-04-28/02 等 - 01 和 02 是子目录,表示数据流流处理管道处理数据的时间。
例子:
gs://data/2020-04-28/01/0000000.avro
gs://data/2020-04-28/01/0000001.avro
gs://data/2020-04-28/01/....
gs://data/2020-04-28/02/0000000.avro
gs://data/2020-04-28/02/0000001.avro
gs://data/2020-04-28/02/....
gs://data/2020-04-28/03/0000000.avro
gs://data/2020-04-28/03/0000001.avro
gs://data/2020-04-28/03/....
...
0000000、0000001 等是我用来说明的简单文件名,我不希望这些文件是顺序名称。 您认为这在 GCP 数据流流设置中可行吗?
【问题讨论】:
标签: google-cloud-platform pipeline dataflow apache-beam