【问题标题】:Pipeline not working in Apache Airflow with Pandas管道在 Pandas 的 Apache Airflow 中不起作用
【发布时间】:2020-12-09 11:34:39
【问题描述】:

我正在尝试形成一个简单的管道,我首先清理数据框,然后将其存储在 CSV 中。 clean_Data 任务失败,因此任务 store_data 不工作。我在store_data 任务中使用完整路径将数据帧写入为CSV。我还首先尝试在本地读取 clean_Data 中的 csv,然后尝试在不从 CSV 读取的情况下创建数据框,但两者都不起作用。

这是我写的 DAG 脚本:

from airflow import DAG
import os
from datetime import timedelta
from datetime import datetime
from datetime import date

from airflow.operators.bash_operator import BashOperator
from airflow.operators.python_operator import PythonOperator
from airflow.utils.dates import days_ago
# import pandas and numpy
import pandas as pd
import numpy as np


default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2020, 12, 5),
    'retries': 1
}


def storeData(**context):
    df= context['task_instance'].xcom_pull(task_ids='clean_Data')
    df.to_csv('/c/Users/HP/airflowhome/dags/data/namesAge2.csv')


def cleanData():
    #df = df.read_csv('/c/Users/HP/airflowhome/dags/data/namesAge.csv')
    data = {'Name': ['Tom', 'nick', 'krish', 'jack'],
            'Age': [np.nan, 21, np.nan, 18]}
    df = pd.DataFrame(data)
    df = df.fillna(0)
    return df


dag = DAG(
    'CleaningPipelineDAG',
    default_args=default_args,
    description='Cleaning Data',
    schedule_interval='@once',
)


t1 = PythonOperator(
    task_id='clean_Data',
    provide_context=True,
    python_callable=cleanData,
    dag=dag,
)

t2 = PythonOperator(
    task_id='store_data',
    provide_context=True,
    python_callable=storeData,
    dag=dag,
)

t1 >> t2

这是日志文件的内容:

[2020-12-09 13:50:55,984] {taskinstance.py:670} INFO - Dependencies all met for <TaskInstance: CleaningPipelineDAG.clean_Data 2020-12-09T11:45:26.969635+00:00 [queued]>
[2020-12-09 13:50:55,992] {taskinstance.py:670} INFO - Dependencies all met for <TaskInstance: CleaningPipelineDAG.clean_Data 2020-12-09T11:45:26.969635+00:00 [queued]>
[2020-12-09 13:50:55,992] {taskinstance.py:880} INFO - 
--------------------------------------------------------------------------------
[2020-12-09 13:50:55,992] {taskinstance.py:881} INFO - Starting attempt 2 of 2
[2020-12-09 13:50:55,992] {taskinstance.py:882} INFO - 
--------------------------------------------------------------------------------
[2020-12-09 13:50:56,004] {taskinstance.py:901} INFO - Executing <Task(PythonOperator): clean_Data> on 2020-12-09T11:45:26.969635+00:00
[2020-12-09 13:50:56,020] {standard_task_runner.py:54} INFO - Started process 1698 to run task
[2020-12-09 13:50:56,091] {standard_task_runner.py:77} INFO - Running: ['airflow', 'run', 'CleaningPipelineDAG', 'clean_Data', '2020-12-09T11:45:26.969635+00:00', '--job_id', '76', '--pool', 'default_pool', '--raw', '-sd', 'DAGS_FOLDER//c/Users/HP/DAG-3.py', '--cfg_path', '/tmp/tmpj92yjpoi']
[2020-12-09 13:50:56,094] {standard_task_runner.py:78} INFO - Job 76: Subtask clean_Data
[2020-12-09 13:50:56,169] {logging_mixin.py:112} INFO - Running <TaskInstance: CleaningPipelineDAG.clean_Data 2020-12-09T11:45:26.969635+00:00 [running]> on host mg.localdomain
[2020-12-09 13:50:56,189] {taskinstance.py:1150} ERROR - cleanData() got an unexpected keyword argument 'conf'
Traceback (most recent call last):
  File "/home/mg/.local/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 984, in _run_raw_task
    result = task_copy.execute(context=context)
  File "/home/mg/.local/lib/python3.8/site-packages/airflow/operators/python_operator.py", line 113, in execute
    return_value = self.execute_callable()
  File "/home/mg/.local/lib/python3.8/site-packages/airflow/operators/python_operator.py", line 118, in execute_callable
    return self.python_callable(*self.op_args, **self.op_kwargs)
TypeError: cleanData() got an unexpected keyword argument 'conf'
[2020-12-09 13:50:56,191] {taskinstance.py:1187} INFO - Marking task as FAILED. dag_id=CleaningPipelineDAG, task_id=clean_Data, execution_date=20201209T114526, start_date=20201209T115055, end_date=20201209T115056
[2020-12-09 13:51:00,984] {local_task_job.py:102} INFO - Task exited with return code 1

【问题讨论】:

  • 究竟是什么不工作?例外?有什么问题?
  • @Elad 我已经用日志文件更新了我的问题。在 Airflow UI 本身上,clean_Data 任务的状态为失败。

标签: python-3.x pandas airflow


【解决方案1】:

您在 cleanData 可调用对象中缺少 **context

这段代码执行正常:

import pandas as pd
import numpy as np

from airflow import DAG
from datetime import datetime
from airflow.operators.python_operator import PythonOperator


default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2020, 12, 5),
    'retries': 1
}


def storeData(**context):
    df= context['task_instance'].xcom_pull(task_ids='clean_Data')
    print (df)


def cleanData(**context):
    data = {'Name': ['Tom', 'nick', 'krish', 'jack'],
            'Age': [np.nan, 21, np.nan, 18]}
    df = pd.DataFrame(data)
    df = df.fillna(0)
    return df


dag = DAG(
    'CleaningPipelineDAG',
    default_args=default_args,
    description='Cleaning Data',
    schedule_interval='@once',
)


t1 = PythonOperator(
    task_id='clean_Data',
    provide_context=True,
    python_callable=cleanData,
    dag=dag,
)

t2 = PythonOperator(
    task_id='store_data',
    provide_context=True,
    python_callable=storeData,
    dag=dag,
)

t1 >> t2

输出:

【讨论】:

  • 你为什么要在 cleanData 中添加 **context,即使我不会获取外部参数并将它们插入到函数中?
  • @MarwaA。然后你需要删除provide_context=True :)
  • 我在之前的 stackoveflow 回答中读到,生成参数的函数和接受参数的函数都应该有 provide_context=True 是不是正确的?
  • 查看 PythonOperator 文档:provide_context: if set to true, Airflow will pass a set of keyword arguments that can be used in your function. This set of kwargs correspond exactly to what you can use in your jinja templates. For this to work, you need to define **kwargs in your function header.
猜你喜欢
  • 2016-08-03
  • 2023-02-08
  • 1970-01-01
  • 2018-02-02
  • 2017-05-28
  • 1970-01-01
  • 2018-01-25
  • 2017-07-02
  • 2019-10-13
相关资源
最近更新 更多