You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 08:54:10