【发布时间】:2020-01-17 02:04:47
【问题描述】:
我有一个程序,它在 pubSub 中创建一个主题并将消息发布到该主题。我还有一个自动数据流作业(使用模板),它将这些消息保存到我的 BigQuery 表中。现在我打算用 python 管道替换基于模板的作业,我的要求是从 PubSub 读取数据,应用转换并将数据保存到 BigQuery/发布到另一个 PubSub 主题。我开始用 python 编写脚本并做了很多试验和错误来实现它,但令我沮丧的是,我无法实现它。代码如下所示:
import apache_beam as beam
from apache_beam.io import WriteToText
TOPIC_PATH = "projects/test-pipeline-253103/topics/test-pipeline-topic"
OUTPUT_PATH = "projects/test-pipeline-253103/topics/topic-repub"
def run():
o = beam.options.pipeline_options.PipelineOptions()
p = beam.Pipeline(options=o)
print("I reached here")
# # Read from PubSub into a PCollection.
data = (
p
| "Read From Pub/Sub" >> beam.io.ReadFromPubSub(topic=TOPIC_PATH)
)
data | beam.io.WriteToPubSub(topic=OUTPUT_PATH)
print("Lines: ", data)
run()
如果我能尽早得到一些帮助,我将不胜感激。 注意:我在谷歌云上设置了我的项目,并且我的脚本在本地运行。
【问题讨论】:
标签: google-cloud-platform google-bigquery google-cloud-dataflow apache-beam google-cloud-pubsub