【发布时间】:2017-07-31 09:03:56
【问题描述】:
以下是应该从 csv 文件读取并写入另一个 csv 文件和 BigQuery 的代码:
import argparse
import logging
import re
import apache_beam as beam
from apache_beam.io import ReadFromText
from apache_beam.io import WriteToText
from apache_beam.metrics import Metrics
from apache_beam.metrics.metric import MetricsFilter
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.options.pipeline_options import SetupOptions
parser = argparse.ArgumentParser()
parser.add_argument('--input',
dest='input',
default='gs://dataflow-samples/shakespeare/kinglear.txt',
help='Input file to process.')
parser.add_argument('--output',
dest='output',
required=True,
help='Output file to write results to.')
known_args, pipeline_args = parser.parse_known_args(None)
pipeline_options = PipelineOptions(pipeline_args)
pipeline_options.view_as(SetupOptions).save_main_session = True
p = beam.Pipeline(options=pipeline_options)
# Read the text file[pattern] into a PCollection.
lines = p | 'read' >> ReadFromText(known_args.input)
lines | beam.Map(lambda x: x.split(','))
lines | 'write' >> WriteToText(known_args.output)
lines | 'write2' >> beam.io.Write(beam.io.BigQuerySink('xxxx:yyyy.aaaa'))
# Actually run the pipeline (all operations above are deferred).
result = p.run()
它能够写入输出文件,但无法写入 BigQuery 表 (xxxx:yyyy.aaaa)
以下是出现的消息:
WARNING:root:A task failed with exception.
'unicode' object has no attribute 'iteritems'
即使架构相同且 BigQuery 表为空,也不会将 csv 文件中包含的表写入 BigQuery。我怀疑这是因为必须将数据转换为 JSON 格式。 为了使其正常工作,必须对此代码进行哪些更正?您能否提供我必须添加的代码行才能使其正常工作?
【问题讨论】: