【发布时间】:2021-01-29 17:37:39
【问题描述】:
我一直在考虑如何在 Beam 中解决给定问题,并认为我会向更多的受众寻求建议。目前事情似乎很少工作,我很好奇是否有人可以提供一个共鸣板来看看这个工作流程是否有意义。
主要的高级目标是从 Kafka 中读取可能出现故障并需要根据记录上的另一个属性在 Event Time 中窗口化的记录,并最终发出这些窗口的内容并将它们写入GCS。
当前管道大致如下所示:
val partitionedEvents = pipeline
.apply("Read Events from Kafka",
KafkaIO
.read<String, Log>()
.withBootstrapServers(options.brokerUrl)
.withTopic(options.incomingEventsTopic)
.withKeyDeserializer(StringDeserializer::class.java)
.withValueDeserializerAndCoder(
SpecificAvroDeserializer<Log>()::class.java,
AvroCoder.of(Log::class.java)
)
.withReadCommitted()
.commitOffsetsInFinalize()
// Set the watermark to use a specific field for event time
.withTimestampPolicyFactory { _, previousWatermark -> WatermarkPolicy(previousWatermark) }
.withConsumerConfigUpdates(
ImmutableMap.of<String, Any?>(
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest",
ConsumerConfig.GROUP_ID_CONFIG, "log-processor-pipeline",
"schema.registry.url", options.schemaRegistryUrl
)
).withoutMetadata()
)
.apply("Logging Incoming Logs", ParDo.of(Events.log()))
.apply("Rekey Logs by Tenant", ParDo.of(Events.key()))
.apply("Partition Logs by Source",
// This is a custom function that will partition incoming records by a specific
// datasource field
Partition.of(dataSources.size, Events.partition<KV<String, Log>>(dataSources))
)
dataSources.forEach { dataSource ->
// Store a reference to the data source name to avoid serialization issues
val sourceName = dataSource.name
val tempDirectory = Directories.resolveTemporaryDirectory(options.output)
// Grab all of the events for this specific partition and apply the source-specific windowing
// strategies
partitionedEvents[dataSource.partition]
.apply(
"Building Windows for $sourceName",
SourceSpecificWindow.of<KV<String, Log>>(dataSource)
)
.apply("Group Windowed Logs by Key for $sourceName", GroupByKey.create())
.apply("Log Events After Windowing for $sourceName", ParDo.of(Events.logAfterWindowing()))
.apply(
"Writing Windowed Logs to Files for $sourceName",
FileIO.writeDynamic<String, KV<String, MutableIterable<Log>>>()
.withNumShards(1)
.by { row -> "${row.key}/${sourceName}" }
.withDestinationCoder(StringUtf8Coder.of())
.via(Contextful.fn(SerializableFunction { logs -> Files.stringify(logs.value) }), TextIO.sink())
.to(options.output)
.withNaming { partition -> Files.name(partition)}
.withTempDirectory(tempDirectory)
)
}
以更简单的项目符号形式,它可能如下所示:
- 从单个 Kafka 主题读取记录
- 按租户键入所有记录
- 由另一个事件正确划分流
- 在上一步中遍历已知分区
- 为每个分区应用自定义窗口规则(与数据源、自定义窗口规则相关)
- 按键分组窗口项(租户)
- 通过 FileIO 将租户密钥对分组写入 GCP
问题是传入的 Kafka 主题包含跨多个租户的乱序数据(例如,租户 1 的事件现在可能正在流入,但几分钟后,您将在相同的分区等)。这将导致水印及时来回反弹,因为不能保证每条传入的记录都会不断增加,这听起来会是一个问题,但我不确定。显然,当数据流过时,某些文件根本没有发出。
自定义窗口函数非常简单,旨在在允许的延迟和窗口持续时间过去后发出单个窗口:
object SourceSpecificWindow {
fun <T> of(dataSource: DataSource): Window<T> {
return Window.into<T>(FixedWindows.of(dataSource.windowDuration()))
.triggering(Never.ever())
.withAllowedLateness(dataSource.allowedLateness(), Window.ClosingBehavior.FIRE_ALWAYS)
.discardingFiredPanes()
}
}
但是,这似乎不一致,因为我们会在窗口关闭后看到日志记录,但不一定是文件被写入 GCS。
这种方法有什么明显错误或不正确的地方吗?由于数据可能在源中出现乱序(即现在、2 小时前、5 分钟后)并涵盖多个租户的数据,但目的是尝试确保一个保持最新状态的租户获胜不会淹没过去可能来的租户。
我们是否可能需要另一个 Beam 应用程序或其他东西来将这个单一的事件流“拆分”成子流,每个子流都独立处理(以便每个水印自行处理)?那是SplittableDoFn 会出现的地方吗?由于我在 SparkRunner 上运行,它似乎不支持这一点 - 但它似乎是一个有效的用例。
任何建议都将不胜感激,甚至只是另一双眼睛。我很乐意提供我能提供的任何其他详细信息。
环境
- 目前针对 SparkRunner 运行
【问题讨论】:
标签: apache-spark apache-kafka triggers apache-beam multi-tenant