Apache Beam/Dataflow中beam.io.ReadFromPubSub输出PCollection的定义及元素归属
Apache Beam/Dataflow中ReadFromPubSub的输出PCollection定义及规则
输出PCollection的类型定义
beam.io.ReadFromPubSub的输出是一个单一的PCollection,其元素类型取决于配置:
- 默认配置下,元素为
PubsubMessage对象,包含消息体(data字段,字节类型)、消息属性(attributes字典)等元数据 - 若指定
with_attributes=False参数,元素则为原始的字节串(bytes类型)
流式消息的归属规则
当依次流式传输10条PubSub消息时,所有消息都会归属到同一个PCollection中,不会各自形成仅含单个元素的PCollection。
PCollection是Beam中用于表示数据集(包括批量、流式数据)的核心抽象,在流式场景下它属于「无界PCollection」,会随着新消息的持续到来不断向集合中添加元素,每条传入的消息都是这个集合里的一个独立元素。
内容的提问来源于stack exchange,提问作者Pav3k
相关产品推荐
相关产品推荐

