如何在Python编写的Airflow DAG中调用Spark Scala函数查询Hive日期
方案1:直接在Python DAG中调用Scala函数的实现方式
你无法在Python的ShortCircuitOperator回调函数中直接调用Scala函数,二者运行环境完全隔离,需要通过中间层调用:
- 第一步:修正你现有Scala函数的语法问题(入参声明了
df但内部又重新定义了同名变量),将其封装为可独立执行的Spark作业,最终把查询到的最新date_of_job_run打印到标准输出,再将整个Spark程序打包为可提交的jar包,指定主类路径。 - 第二步:在Python的
getLatestRunDate函数中,通过subprocess调用spark-submit命令执行jar包,解析返回的日期结果,和DAG执行日期做对比。
对应Python代码示例:
import subprocess from datetime import datetime def getLatestRunDate(**context): # 调用打包好的Scala Spark作业 cmd = "spark-submit --class com.yourcompany.utils.GetLatestRunDate /path/to/your/scala-job.jar" run_date = subprocess.check_output(cmd, shell=True).decode("utf-8").strip() # 取DAG执行日期和查询结果对比 exec_date = context["ds"] return run_date == exec_date
该方案每次执行需要拉起Spark进程,资源开销较大,仅适合必须复用复杂Scala逻辑的场景。
方案2:更优的推荐实现方案
你的需求仅为查询Hive表中的日期做匹配,完全不需要绕Scala调用,直接在Python侧轻量查询Hive即可,性能和维护成本都远好于方案1:
- 方式1:用PyHive连接HiveServer2直接查询,无Spark额外开销
from pyhive import hive def getLatestRunDate(**context): conn = hive.connect(host='你的HiveServer2地址', port=10000, username='airflow') cursor = conn.cursor() # 直接把Scala中的过滤逻辑转成SQL即可 query_sql = """ SELECT date_of_job_run FROM my_hive_schema.my_run_catalog_table WHERE date_of_job_run <= current_date() AND month_nbr >= month(current_date()) - 2 AND yr_nbr = year(current_date()) ORDER BY date_of_job_run DESC LIMIT 1 """ cursor.execute(query_sql) latest_run_date = str(cursor.fetchone()[0]) exec_date = context["ds"] return latest_run_date == exec_date
- 方式2:如果Airflow已经配置了Hive连接,直接用官方Hive Hook更规范,不需要自己管理连接
from airflow.providers.apache.hive.hooks.hive import HiveServer2Hook def getLatestRunDate(**context): hive_hook = HiveServer2Hook(hive_conn_id='你配置的Hive连接ID') query_sql = """ SELECT date_of_job_run FROM my_hive_schema.my_run_catalog_table WHERE date_of_job_run <= current_date() AND month_nbr >= month(current_date()) - 2 AND yr_nbr = year(current_date()) ORDER BY date_of_job_run DESC LIMIT 1 """ latest_run_date = str(hive_hook.get_records(query_sql)[0][0]) exec_date = context["ds"] return latest_run_date == exec_date
内容的提问来源于stack exchange,提问作者palamuGuy
相关产品推荐
相关产品推荐

