【问题标题】:Airflow Task Succeeded But Not All Data IngestedAirflow 任务成功但未摄取所有数据
【发布时间】:2023-02-14 16:16:43
【问题描述】:

我有一个气流任务来用这个流提取数据

PostgreSQL -> Google Cloud Storage -> BigQuery

我遇到的问题是,似乎并非所有数据都被提取到 BigQuery 中。在 PostgreSQL 源上,该表有 18M+ 行数据,但在摄取后它只有 4M+ 行数据。

当我检查生产时,数据通过此查询返回 18M+ 行:

SELECT COUNT(1) FROM my_table

-- This return 18M+ rows

但是在 DAG 完成运行后,当我检查 BigQuery 时:

SELECT COUNT(1) FROM data_lake.my_table

-- This return 4M+ rows

为了做笔记,并非我摄取的所有表格都是这样返回的。所有较小的桌子都摄取得很好。但是当它达到一定数量的行时,它的行为就像这样。

我怀疑数据是从 PostgreSQL 提取到 Google Cloud Storage 的。所以我会在这里提供我的功能:

    def create_operator_write_append_init(self, worker=10):
        worker_var  = dict()
        with TaskGroup(group_id=self.task_id_init) as tg1:
            for i in range(worker):
                worker_var[f'worker_{i}'] = PostgresToGCSOperator(
                    task_id = f'worker_{i}',
                    postgres_conn_id = self.conn_id,
                    sql = 'extract_init.sql',
                    bucket = self.bucket,
                    filename = f'{self.filename_init}_{i}.{self.export_format}',
                    export_format = self.export_format, # the export format is json
                    gzip = True,
                    params = {
                        'worker': i
                    }
                )
        return tg1

这是 SQL 文件:

SELECT id,
       name,
       created_at,
       updated_at,
       deleted_at
FROM my_table
WHERE 1=1
AND ABS(MOD(hashtext(id::TEXT), 10)) = {{params.worker}};

我所做的是将数据分块并将其分成几个工作人员,因此是 TaskGroup。

提供更多信息。我使用作曲家:

  • 作曲家-2.0.32-气流-2.3.4

  • 大型实例

  • 工人 8CPU

  • 工人 32GB 内存

  • 工人 2GB 存储空间

  • 1-16 岁之间的工人

这些发生的可能性有多大?

【问题讨论】:

    标签: python airflow google-cloud-composer airflow-2.x


    【解决方案1】:

    PostgresToGCSOperator继承自BaseSQLToGCSOperator(https://airflow.apache.org/docs/apache-airflow-providers-google/stable/_api/airflow/providers/google/cloud/transfers/sql_to_gcs/index.html)

    根据源代码,approx_max_file_size_bytes=1900000000。所以如果你把你的表分成 10 个部分(或者工人说)这个块的最大大小应该是最大 1.9 GB。如果这个块更大,以前的块将被新块替换,因为你没有指定由 PostgresToGCSOperator 创建“你的块的块”。

    您可以通过在 filename 中添加占位符 {} 来实现它,Operator 将处理它。

    def create_operator_write_append_init(self, worker=10):
            worker_var  = dict()
            with TaskGroup(group_id=self.task_id_init) as tg1:
                for i in range(worker):
                    worker_var[f'worker_{i}'] = PostgresToGCSOperator(
                        task_id = f'worker_{i}',
                        postgres_conn_id = self.conn_id,
                        sql = 'extract_init.sql',
                        bucket = self.bucket,
                        filename = f'{self.filename_init}_{i}_part_{{}}.{self.export_format}',                    
                        export_format = self.export_format, # the export format is json
                        gzip = True,
                        params = {
                            'worker': i
                        }
                    )
            return tg1
    

    【讨论】:

    • 感谢你的回答!我会尝试占位符并回复您。因为上周我需要这么快,所以我所做的是垂直缩放作曲家实例。它有效,但我知道这不是最佳做法。如果还有其他需求需要我再次摄取大表,我会尝试这种方法,如果可行的话会回来。
    • 如果这有助于您理解一个概念,我能否请您对我的回答投赞成票?
    • 是的,我已经对你的回答投了赞成票,但我还没有足够的声誉,因为它可以在帖子中看到:)
    【解决方案2】:

    您绝对可以探索由 Astronomer 维护的 Apache 2.0 许可Astro SDK,它允许使用由 Apache Airflow 提供支持的 Python 和 SQL 快速干净地开发{Extract,Load,Transform}工作流。

    在这种情况下,aql.transform_file 可用于从 .sql 文件运行 SQL 查询并从 Postgres 中选择数据。 aql.export_to_file() 会将数据从 Postgres 表导出到 GCS 存储桶。最后,aql.load_file() 可用于将文件中的数据从 GCS 加载到 BigQuery。以下是示例 DAG:

    from airflow.models.dag import DAG
    
    from astro.files import File
    from astro.constants import FileType
    from astro.table import Table
    from astro.sql.operators.load_file import load_file
    from astro.sql.operators.export_to_file import export_to_file
    from astro.sql.operators.transform import transform_file
    from datetime import datetime
    from pathlib import Path
    
    POSTGRES_CONN_ID ="postgres_conn"
    
    with DAG(
            dag_id="sample-dag",
            schedule_interval=None,
            start_date=datetime(2022, 1, 1),
            catchup=False,
    ) as dag:
        postgres_table = Table(name="my_table", temp=True, conn_id=POSTGRES_CONN_ID)
    
        postgres_data = transform_file(
            file_path=f"{Path(__file__).parent.as_posix()}/transform.sql",
            parameters={"input_table": postgres_table},
        )
    
    
        save_file_to_gcs = export_to_file(
            task_id="save_file_to_gcs",
            input_data=postgres_data,
            output_file=File(
                path="gs://astro-sdk/all_postgres_data.csv",
                conn_id="gcp_conn",
            ),
            if_exists="replace",
        )
    
        load_data_to_bq = load_file(
            input_file=File(
                "gs://astro-sdk/all_postgres_data.csv",
                conn_id="gcp_conn",
                filetype=FileType.CSV,
            ),
            output_table=Table(conn_id="gcp_conn"),
            use_native_support=False,
            native_support_kwargs={
                "ignore_unknown_values": True,
                "allow_jagged_rows": True,
                "skip_leading_rows": "1",
            },
            enable_native_fallback=True,
        )
        load_data_to_bq.set_upstream(save_file_to_gcs)
    

    添加 DAG 运行的屏幕截图。 DAG screenshot

    因此,改用 astro-sdk-python 只会简化方法。

    我们有各种操作员和装饰器作为这个项目的一部分,在这里描述:https://astro-sdk-python.readthedocs.io/

    免责声明:我在Astronomer 工作,该公司将 Astro SDK 开发为开源项目。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-07-21
      • 1970-01-01
      • 2021-02-22
      • 2021-12-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多