将消息从 pubsubio 插入 BigQuery 时出现错误提示。
如何将记录从 pubsub 插入 BQ。我们可以转换pcollection
成一个列表,还是有另一个替代品?
AttributeError:
'PCollection'
对象没有属性'split'
这是我的代码:
def create_record(columns):
#import re
col_value=record_ids.split('|')
col_name=columns.split(",")
for i in range(length(col_name)):
schmea_dict[col_name[i]]=col_value[i]
return schmea_dict
schema = 'tungsten_opcode:STRING,tungsten_seqno:INTEGER
columns="tungsten_opcode,tungsten_seqno"
lines = p | 'Read PubSub' >> beam.io.ReadStringsFromPubSub(INPUT_TOPIC) |
beam.WindowInto(window.FixedWindows(15))
record_ids = lines | 'Split' >>
(beam.FlatMap(split_fn).with_output_types(unicode))
records = record_ids | 'CreateRecords' >> beam.Map(create_record(columns))
records | 'BqInsert' >> beam.io.WriteToBigQuery(
OUTPUT,
schema=schema,
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND)