【问题标题】:Partitioning a table对表进行分区
【发布时间】:2017-10-13 13:26:32
【问题描述】:

Bigquery 目前只允许按日期进行分区。

假设我有一个包含 inserted_timestamp 字段的 10 亿表行。假设该字段的日期为 1 年前。

将现有数据移动到新分区表的正确方法是什么?

已编辑

我看到有一个关于 Java 版本 Sharding BigQuery output tables 的优雅解决方案也在 BigQuery partitioning with Beam streams 详细说明,即参数化表名(或分区后缀)窗口数据。

但我在 2.x 梁项目上想念BigQueryIO.Write,也没有关于从 python 可序列化函数获取窗口时间的示例。

我尝试在管道上创建分区,但如果因大量分区而失败(以 100 运行但因 1000 失败)。

这是我的代码:

               (  p
                | 'lectura' >> beam.io.ReadFromText(input_table)
                | 'noheaders' >> beam.Filter(lambda s: s[0].isdigit())
                | 'addtimestamp' >> beam.ParDo(AddTimestampDoFn())
                | 'window' >> beam.WindowInto(beam.window.FixedWindows(60))
                | 'table2row'  >> beam.Map( to_table_row )  
                | 'write2table' >> beam.io.Write(beam.io.BigQuerySink(
                        output_table,   #<-- unable to parametrize by window
                        dataset=my_dataset, 
                        project=project, 
                        schema='dia:DATE, classe:STRING, cp:STRING, import:FLOAT',
                        create_disposition=CREATE_IF_NEEDED,
                        write_disposition=WRITE_TRUNCATE,
                                    )
                                )
                )

p.run()

【问题讨论】:

  • stackoverflow.com/questions/38993877/… 应该有一些相关的方法。另外我认为您应该能够使用 JSON 或 AVRO 而不是 CSV 来避免使用平面文件。
  • @NhanNguyen,刚刚将我的问题编辑得更具体。在 2.x 上想念它。感谢您的链接,我关注了它并且是非常相关的问题。再次感谢。

标签: python google-bigquery google-cloud-dataflow


【解决方案1】:

所有必要的功能都存在于 Beam 中,但目前可能仅限于 Java SDK。

您将使用BigQueryIO。具体来说,您可以使用DynamicDestinations 来确定每一行的目标表。

以 DynamicDestinations 为例:

events.apply(BigQueryIO.<UserEvent>write()
  .to(new DynamicDestinations<UserEvent, String>() {
        public String getDestination(ValueInSingleWindow<String> element) {
          return element.getValue().getUserId();
        }
        public TableDestination getTable(String user) {
          return new TableDestination(tableForUser(user), 
            "Table for user " + user);
        }
        public TableSchema getSchema(String user) {
          return tableSchemaForUser(user);
        }
      })
  .withFormatFunction(new SerializableFunction<UserEvent, TableRow>() {
     public TableRow apply(UserEvent event) {
       return convertUserEventToTableRow(event);
     }
   }));

【讨论】:

  • 为什么它们不是一个 python 包装器来做呢?我应该用Java而不是python来负担数据流项目吗?你知道 Google 是否将资源集中在 Java 上吗?我的意思是,如果我在 Python 中工作,我会错过比这个更多的功能吗?谢谢!
  • 如上所示,Java 和 Python SDK 之间存在不同的功能。解决这些差距是 Apache Beam 持续努力的一部分。此特定问题被跟踪为BEAM-2801
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-05-05
  • 2018-12-16
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-02-16
相关资源
最近更新 更多