【问题标题】:Google Cloud Dataflow Write to CSV from dictionaryGoogle Cloud Dataflow 从字典写入 CSV
【发布时间】:2017-10-17 15:57:30
【问题描述】:

我有一个值字典,我想使用 Python SDK 将其作为有效的 .CSV 文件写入 GCS。我可以将字典写成换行符分隔的文本文件,但我似乎找不到将字典转换为有效 .CSV 的示例。任何人都可以建议在数据流管道中生成 csv 的最佳方法吗?这回答了这个question 地址读取 CSV 文件,但并没有真正解决写入 CSV 文件的问题。我认识到 CSV 文件只是带有规则的文本文件,但我仍在努力将数据字典转换为可以使用 WriteToText 编写的 CSV。

这是一个简单的示例字典,我想将其转换为 CSV:

test_input = [{'label': 1, 'text': 'Here is a sentence'},
              {'label': 2, 'text': 'Another sentence goes here'}]


test_input  | beam.io.WriteToText(path_to_gcs)

以上将产生一个文本文件,其中每个字典都在换行符上。 Apache Beam 中是否有任何我可以利用的功能(类似于csv.DictWriter)?

【问题讨论】:

    标签: python google-cloud-dataflow apache-beam


    【解决方案1】:

    通常,您需要编写一个函数,将您的原始 dict 数据元素转换为 csv 格式的 string 表示形式。

    该函数可以编写为DoFn,您可以将其应用于您的数据束PCollection,这会将每个集合元素转换为所需的格式;您可以通过ParDoDoFn 应用到您的PCollection 来做到这一点。您还可以将这个DoFn 包装在一个更用户友好的PTransform 中。

    你可以在Beam Programming Guide了解更多关于这个过程的信息

    这是一个简单的、可翻译的非 Beam 示例:

    # Our example list of dictionary elements
    test_input = [{'label': 1, 'text': 'Here is a sentence'},
                 {'label': 2, 'text': 'Another sentence goes here'}]
    
    def convert_my_dict_to_csv_record(input_dict):
        """ Turns dictionary values into a comma-separated value formatted string """
        return ','.join(map(str, input_dict.values()))
    
    # Our converted list of elements
    converted_test_input = [convert_my_dict_to_csv_record(element) for element in test_input]
    

    converted_test_input 将如下所示:

    ['Here is a sentence,1', 'Another sentence goes here,2']
    

    Beam DictToCSV DoFn 和 PTransform 示例使用DictWriter

    from csv import DictWriter
    from csv import excel
    from cStringIO import StringIO
    
    ...
    
    def _dict_to_csv(element, column_order, missing_val='', discard_extras=True, dialect=excel):
        """ Additional properties for delimiters, escape chars, etc via an instance of csv.Dialect
            Note: This implementation does not support unicode
        """
    
        buf = StringIO()
    
        writer = DictWriter(buf,
                            fieldnames=column_order,
                            restval=missing_val,
                            extrasaction=('ignore' if discard_extras else 'raise'),
                            dialect=dialect)
        writer.writerow(element)
    
        return buf.getvalue().rstrip(dialect.lineterminator)
    
    
    class _DictToCSVFn(DoFn):
        """ Converts a Dictionary to a CSV-formatted String
    
            column_order: A tuple or list specifying the name of fields to be formatted as csv, in order
            missing_val: The value to be written when a named field from `column_order` is not found in the input element
            discard_extras: (bool) Behavior when additional fields are found in the dictionary input element
            dialect: Delimiters, escape-characters, etc can be controlled by providing an instance of csv.Dialect
    
        """
    
        def __init__(self, column_order, missing_val='', discard_extras=True, dialect=excel):
            self._column_order = column_order
            self._missing_val = missing_val
            self._discard_extras = discard_extras
            self._dialect = dialect
    
        def process(self, element, *args, **kwargs):
            result = _dict_to_csv(element,
                                  column_order=self._column_order,
                                  missing_val=self._missing_val,
                                  discard_extras=self._discard_extras,
                                  dialect=self._dialect)
    
            return [result,]
    
    class DictToCSV(PTransform):
        """ Transforms a PCollection of Dictionaries to a PCollection of CSV-formatted Strings
    
            column_order: A tuple or list specifying the name of fields to be formatted as csv, in order
            missing_val: The value to be written when a named field from `column_order` is not found in an input element
            discard_extras: (bool) Behavior when additional fields are found in the dictionary input element
            dialect: Delimiters, escape-characters, etc can be controlled by providing an instance of csv.Dialect
    
        """
    
        def __init__(self, column_order, missing_val='', discard_extras=True, dialect=excel):
            self._column_order = column_order
            self._missing_val = missing_val
            self._discard_extras = discard_extras
            self._dialect = dialect
    
        def expand(self, pcoll):
            return pcoll | ParDo(_DictToCSVFn(column_order=self._column_order,
                                              missing_val=self._missing_val,
                                              discard_extras=self._discard_extras,
                                              dialect=self._dialect)
                                 )
    

    要使用该示例,您可以将test_input 放入PCollection,并将DictToCSV PTransform 应用于PCollection;您可以获取转换后的PCollection 并将其用作WriteToText 的输入。请注意,您必须通过 column_order 参数提供与字典输入元素的键相对应的列名列表或元组;生成的 CSV 格式的字符串列将按照提供的列名的顺序排列。此外,该示例的底层实现不支持unicode

    【讨论】:

    • 谢谢,安德鲁!这对我来说很有意义 - 认为我了解机制,我想知道的是 Apache-Beam 中是否存在类似 ConverDictToCSVFn() 的东西,或者是否必须从头开始编写。编写这种类型的函数并非易事,因为如果句子包含逗号(或任何分隔符),那么您通常需要用双引号“”将整个句子括起来。我猜这个响应表明 Apache-Beam 中没有任何东西可以处理这些情况?
    • 据我所知,textio 似乎没有这种便利性——同时,我相信这可以通过将csv 模块的DictWriter 与Python StringIO 模块。
    • 您能否就您的建议提供更多指导?也许有一个单独的答案?这确实是我问题的症结所在,因为我一直无法找到使用 DictWriter 的方法......
    • 当然——我将使用DictWriterStringIO 实现这一点的一些代码更新上面的响应——我将与Beam SDK 团队合作,看看我们是否可以得到这也是通过拉取请求添加的。
    • 谢谢!我根据您之前的建议在下面添加了一个答案,但我还没有弄清楚如何使用 DictWriter ,所以那会很棒。
    【解决方案2】:

    根据 Andrew 的建议,这是我创建的 ConvertDictToCSV 函数:

    def ConvertDictToCSV(input_dict, fieldnames, separator=",", quotechar='"'):
      value_list = []
      for field in fieldnames:
        if input_dict[field]:
          field_value = str(input_dict[field])
        else:
          field_value = ""
        if separator in field_value:
          field_value = quotechar + field_value + quotechar
        value_list.append(field_value)
    
      return separator.join(value_list)
    

    这似乎运作良好,但如果可能的话,使用 csv.DictWriter 肯定会更安全

    【讨论】:

    • 如果field_value 包含quotechar,你会遇到问题,因为这不会被转义。
    • 分布式系统中尽量不要使用for循环
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-11-02
    • 2018-03-19
    相关资源
    最近更新 更多