【发布时间】:2019-08-30 15:11:15
【问题描述】:
我正在使用 Celery Executor 运行 Airflow v1.9.0。我已经为不同的工作人员配置了不同的队列名称,例如 DEV、QA、UAT、PROD。我编写了一个自定义传感器,它轮询源数据库连接和目标数据库连接并运行不同的查询并在触发下游任务之前进行一些检查。对于多个工人来说,这一直运行良好。在其中一名工人中,此传感器给出了一个 AttributeError 问题:
$ airflow test PDI_Incr_20190407_v1 checkCCWatermarkDt 2019-04-09
[2019-04-09 10:02:57,769] {configuration.py:206} WARNING - section/key [celery/celery_ssl_active] not found in config
[2019-04-09 10:02:57,770] {default_celery.py:41} WARNING - Celery Executor will run without SSL
[2019-04-09 10:02:57,771] {__init__.py:45} INFO - Using executor CeleryExecutor
[2019-04-09 10:02:57,817] {models.py:189} INFO - Filling up the DagBag from /home/airflow/airflow/dags
/usr/local/lib/python2.7/site-packages/airflow/models.py:2160: PendingDeprecationWarning: Invalid arguments were passed to ExternalTaskSensor. Support for passing such arguments will be dropped in Airflow 2.0. Invalid arguments were:
*args: ()
**kwargs: {'check_existence': True}
category=PendingDeprecationWarning
[2019-04-09 10:02:57,989] {base_hook.py:80} INFO - Using connection to: 172.16.20.11:1521/GWPROD
[2019-04-09 10:02:57,991] {base_hook.py:80} INFO - Using connection to: dmuat.cwmcwghvymd3.us-east-1.rds.amazonaws.com:1521/DMUAT
Traceback (most recent call last):
File "/usr/local/bin/airflow", line 27, in <module>
args.func(args)
File "/usr/local/lib/python2.7/site-packages/airflow/bin/cli.py", line 528, in test
ti.run(ignore_task_deps=True, ignore_ti_state=True, test_mode=True)
File "/usr/local/lib/python2.7/site-packages/airflow/utils/db.py", line 50, in wrapper
result = func(*args, **kwargs)
File "/usr/local/lib/python2.7/site-packages/airflow/models.py", line 1584, in run
session=session)
File "/usr/local/lib/python2.7/site-packages/airflow/utils/db.py", line 50, in wrapper
result = func(*args, **kwargs)
File "/usr/local/lib/python2.7/site-packages/airflow/models.py", line 1493, in _run_raw_task
result = task_copy.execute(context=context)
File "/usr/local/lib/python2.7/site-packages/airflow/operators/sensors.py", line 78, in execute
while not self.poke(context):
File "/home/airflow/airflow/plugins/PDIPlugin.py", line 29, in poke
wm_dt_src = hook_src.get_records(self.sql)
AttributeError: 'NoneType' object has no attribute 'get_records'
虽然当我从调度程序 CLI 运行相同的测试命令时,它运行良好。上面的问题看起来像一个数据库连接问题。
为了调试,我从 Airflow UI 检查了 DB Connections: 数据分析 -> 即席查询 查询:从对偶中选择1; -- 这很好用
我还从工作节点远程登录到数据库主机和端口,也很顺利。
自定义传感器代码:
from airflow.plugins_manager import AirflowPlugin
from airflow.hooks.base_hook import BaseHook
from airflow.operators.sensors import SqlSensor
class SensorWatermarkDt(SqlSensor):
def __init__(self, conn_id, sql, conn_id_tgt, sql_tgt, *args, **kwargs):
self.sql = sql
self.conn_id = conn_id
self.sql_tgt = sql_tgt
self.conn_id_tgt = conn_id_tgt
super(SqlSensor, self).__init__(*args, **kwargs)
def poke(self, context):
hook_src = BaseHook.get_connection(self.conn_id).get_hook()
hook_tgt = BaseHook.get_connection(self.conn_id_tgt).get_hook()
self.log.info('Poking: %s', self.sql)
self.log.info('Poking: %s', self.sql_tgt)
wm_dt_src = hook_src.get_records(self.sql)
wm_dt_tgt = hook_tgt.get_records(self.sql_tgt)
if wm_dt_src <= wm_dt_tgt:
return False
else:
return True
class PDIPlugin(AirflowPlugin):
name = "PDIPlugin"
operators = [SensorWatermarkDt]
气流 DAG 片段:
import airflow
from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from airflow.operators.email_operator import EmailOperator
from datetime import timedelta,datetime
from airflow.operators import SensorWatermarkDt
from airflow.operators.sensors import ExternalTaskSensor
from airflow.operators.dummy_operator import DummyOperator
default_args = {
'owner': 'SenseTeam',
#'depends_on_past': True,
'depends_on_past' : False,
'start_date': datetime(2019, 4, 7, 17, 00),
'email': [],
'email_on_failure': False,
'email_on_retry': False,
'queue': 'PENTAHO_UAT'
}
dag = DAG(dag_id='PDI_Incr_20190407_v1',
default_args=default_args,
max_active_runs=1,
concurrency=1,
catchup=False,
schedule_interval=timedelta(hours=24),
dagrun_timeout=timedelta(minutes=23*60))
checkCCWatermarkDt = \
SensorWatermarkDt(task_id='checkCCWatermarkDt',
conn_id='CCUSER_SOURCE_GWPROD_RPT',
sql="SELECT MAX(CC_WM.CREATETIME) as CURRENT_WATERMARK_DATE FROM CCUSER.CCX_CAPTUREREASON_ETL CC_WM INNER JOIN CCUSER.CCTL_CAPTUREREASON_ETL CC_WMLKP ON CC_WM.CAPTUREREASON_ETL = CC_WMLKP.ID AND UPPER(CC_WMLKP.DESCRIPTION)= 'WATERMARK'",
conn_id_tgt = 'RDS_DMUAT_DMCONFIG',
sql_tgt = "SELECT MAX(CURRENT_WATERMARK_DATE) FROM DMCONFIG.PRESTG_DM_WMD_WATERMARKDATE WHERE SCHEMA_NAME = 'CCUSER'",
poke_interval=60,
dag=dag)
...
在此工作节点中添加此插件后,我已重新启动 Web 服务器、调度程序和气流工作程序。
我在这里错过了什么?
【问题讨论】:
-
请告诉我们您如何称呼您的
PDIPlugin -
在气流 DAG 片段中添加。但是,我注意到此框中未安装 cx_Oracle。安装和配置 cx_Oracle 后,它工作了。虽然不确定,如果它是正确的方法。如果问题是由于缺少 cx_Oracle 模块而发生的,我希望错误消息更具体。
标签: airflow