【问题标题】:Dataflow GCS to BigQuery - How to output multiple rows per input?数据流 GCS 到 BigQuery - 如何为每个输入输出多行?
【发布时间】:2019-04-10 01:42:55
【问题描述】:

目前我正在使用谷歌提供的 gcs-text-to-bigquery 模板并输入一个转换函数来转换我的 jsonl 文件。 jsonl 是非常嵌套的,我希望能够通过做一些转换来为换行分隔的 json 的每一行输出多行。

例如:

{'state': 'FL', 'metropolitan_counties':[{'name': 'miami dade', 'population':100000}, {'name': 'county2', 'population':100000}…], 'rural_counties':{'name': 'county1', 'population':100000}, {'name': 'county2', 'population':100000}….{}], 'total_state_pop':10000000,….}

显然会有比 2 个更多的县,每个州都会有其中一条线。我老板想要的输出是:

当我进行 gcs-to-bq 文本转换时,我最终每个州只得到一行(所以我将从 FL 获得 miami dade 县,然后无论第一个县在我的转换中为下一个州)。我读了一点,我认为这是因为模板中的映射需要每个 jsonline 一个输出。看来我可以做一个pardo(DoFn?)不确定那是什么,或者在python中有一个与beam.Map类似的选项。转换中有一些业务逻辑(现在大约有 25 行代码,因为 json 的列比我显示的要多,但这些非常简单)。

对此有什么建议吗?今晚/明天会有数据,BQ 表会有几十万行。

我正在使用的模板目前是在 java 中,但我可以很容易地将它翻译成 python,因为在 python 中有很多在线示例。我更了解python,并且我认为它更容易考虑到不同的类型(有时一个字段可以为空),并且鉴于我看到的示例看起来更简单,它似乎不那么令人生畏,但是,对任何一个都开放

【问题讨论】:

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


    【解决方案1】:

    在 Python 中解决这个问题有点简单。这是一种可能性(未完全测试):

    from __future__ import absolute_import                                                               
    
    import ast                                                                      
    
    import apache_beam as beam                                                      
    from apache_beam.io import ReadFromText                                            
    from apache_beam.io import WriteToText                                             
    
    from apache_beam.options.pipeline_options import PipelineOptions                   
    
    import os                                                                       
    os.environ['GOOGLE_APPLICATION_CREDENTIALS'] = '/path/to/service_account.json'      
    
    pipeline_args = [                                                                  
        '--job_name=test'                                                              
    ]                                                                                  
    
    pipeline_options = PipelineOptions(pipeline_args)                                  
    
    
    def jsonify(element):                                                              
        return ast.literal_eval(element)                                               
    
    
    def unnest(element):                                                            
        state = element.get('state')                                                
        state_pop = element.get('total_state_pop')                                  
        if state is None or state_pop is None:                                                   
            return                                                                  
        for type_ in ['metropolitan_counties', 'rural_counties']:                   
            for e in element.get(type_, []):                                        
                name = e.get('name')                                                
                pop = e.get('population')                                           
                county_type = (                                                     
                    'Metropolitan' if type_ == 'metropolitan_counties' else 'Rural' 
                )                                                                   
                if name is None or pop is None:                                     
                    continue                                                        
                yield {                                                             
                    'State': state,                                                 
                    'County_Type': county_type,                                     
                    'County_Name': name,                                            
                    'County_Pop': pop,                                              
                    'State_Pop': state_pop                                          
                }
    
    with beam.Pipeline(options=pipeline_options) as p:                              
        lines = p | ReadFromText('gs://url to file')                                        
    
        schema = 'State:STRING,County_Type:STRING,County_Name:STRING,County_Pop:INTEGER,State_Pop:INTEGER'                                                                      
    
        data = (                                                                    
            lines                                                                   
            | 'Jsonify' >> beam.Map(jsonify)                                        
            | 'Unnest' >> beam.FlatMap(unnest)                                      
            | 'Write to BQ' >> beam.io.Write(beam.io.BigQuerySink(                  
                'project_id:dataset_id.table_name', schema=schema,                     
                create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,    
                write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND)       
            )                                                                       
        )
    

    只有在处理批处理数据时才会成功。如果您有流式数据,只需将beam.io.Write(beam.io.BigquerySink(...)) 更改为beam.io.WriteToBigQuery

    【讨论】:

    • 非常感谢 - 我将如何在 gcs(不是本地)上运行它?我更改了 with beam... 使其成为一个函数,然后有一个 python if name == main 来调用 run 函数以及添加的参数。它在技术上不需要任何输入(目前),但可以不知道如何运行它。
    • 您的意思是运行它从 GCS 检索文件或使用数据流作为运行器?
    • 是的,我在这个线程stackoverflow.com/questions/47821942/…的基础上稍微改变了模板
    • 但我仍然不确定如何运行它(暂时只是一次作为测试)
    • 如果你想在本地测试它,你可以用 python python file_name.py 调用它。在数据流中运行需要将参数 --runner=Dataflow 传递给管道选项。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-04-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-11-06
    • 2011-06-06
    • 1970-01-01
    相关资源
    最近更新 更多