【问题标题】:Is there a way to convert PubSubIO read to UnboundedSource source有没有办法将 PubSubIO 读取转换为 UnboundedSource 源
【发布时间】: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


【解决方案1】:

PubSubIO 目前没有很好地支持这一点,而且对于“Beam 模型”来说有点奇怪。一些选项:

  1. 您是否尝试过启动管道并定期排空它?
  2. 如果这不起作用,您应该在 Beam 邮件列表或issue tracker 中发布功能请求。

【讨论】:

    猜你喜欢
    • 2014-11-06
    • 2017-03-25
    • 2021-09-01
    • 2021-05-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多