【发布时间】:2019-09-06 02:38:25
【问题描述】:
我想使用PubSub subscription 作为有界源,以最大限度地降低流传输管道一直运行的成本。在Batch Pipeline with Unbounded Source 之前提出了类似的问题,但没有解决方案。我遇到了这个答案What PipelineRunners,它说我们可以将UnboundedSource 转换为BoundedSource 以使用withMaxNumRecords 进行测试。是否可以在此处使用PubSubIO 作为输入,或者是否有办法将PubSubIO 读取为unboundedSource?
UnboundedSource<String> unboundedSource = .; // How to Use PubSub here?
PCollection<String> boundedPubsubCollection =
p.apply(Read.from(unboundedSource).withMaxNumRecords(10));
【问题讨论】:
-
您是否尝试使用您的 maxNumRecord 和流参数为 false 将 pubsub 作为无限源读入?
-
是的,我试过了——
UnboundedSource unboundedSource=PubsubIO.readStrings().fromTopic("testTopic");是非法的。无论如何要使用这个? -
是的,当然,这是非法的,它是一个返回的 PCollection。无论如何,你想达到什么目的?你有什么要求?你的目标是什么?通过考虑您的解决方案,我认为数据流不是正确的平台(我不谈论 Beam 编程语言,而只谈论您运行管道的平台。)您能否以更广泛的目标视野来编辑您的问题?
-
@guillaume blaquiere 我的问题与Batch Pipeline with Unbounded Source非常相似。就像将蒸汽转换为批处理管道。任何帮助都感激不尽。谢谢
标签: google-cloud-dataflow apache-beam google-cloud-pubsub