【问题标题】:How to obtain and process mysql records using Airflow?如何使用 Airflow 获取和处理 mysql 记录?
【发布时间】:2018-03-03 17:24:03
【问题描述】:

我需要

1. run a select query on MYSQL DB and fetch the records.              
2. Records are processed by python script.

我不确定我应该如何进行。 xcom是去这里的路吗?此外,MYSQLOperator 只执行查询,不获取记录。我可以使用任何内置的传输运算符吗?如何在这里使用 MYSQL 挂钩?

您可能希望使用 PythonOperator 来使用钩子获取数据, 应用转换并将(现在评分的)行发送回其他地方。

有人可以解释如何处理相同的问题。

参考 - http://markmail.org/message/x6nfeo6zhjfeakfe

def do_work():
    mysqlserver = MySqlHook(connection_id)
    sql = "SELECT * from table where col > 100 "
    row_count = mysqlserver.get_records(sql, schema='testdb')
    print row_count[0][0]

callMYSQLHook = PythonOperator(
    task_id='fetch_from_testdb',
    python_callable=mysqlHook,
    dag=dag
)

这是正确的方法吗? 还有我们如何使用 xcoms 来存储下面的 MySqlOperator 的记录?'

t = MySqlOperator(
conn_id='mysql_default',
task_id='basic_mysql',
sql="SELECT count(*) from table1 where id > 10",
dag=dag)

【问题讨论】:

  • 不要使用 MySqlHook,因为您想要获取大量似乎不适合 XCOM 的数据,您可以使用 PythonOperator 创建自己的函数,然后可能有 sqlalchemy 来帮助您用sql连接查询?
  • @Chengzhi 这就是我目前正在做的,但它违背了当时使用 Airflow 的目的。还有其他解决方法吗?
  • 我不认为这违背了使用气流的目的。事物上的操作符(MySQL 操作符对 MySQL 数据库进行操作)。如果您想使用 Python 对数据库中的每条记录进行操作,那么您只需要使用 PythonOperator 才有意义。我不会害怕制作使用像 sqlalchemy 这样的低级包的大型 Python 脚本。有时,如果这是您需要做的任务,那只是最好的选择。我认为你的问题是你试图从一个真正的任务中做出两个任务。
  • 用Python处理后的记录怎么办?你会把它们写到其他地方吗?有很多社区贡献的操作员,其中许多是多个钩子的组合,在 Pandas 或类似操作之间进行了一些处理,例如以 MySQL 到 GCS 为例。核心运营商目录github.com/apache/incubator-airflow/tree/master/airflow/… 中有很多内容,也不要忘记 contrib 目录中的内容。如果您需要的内容不存在,请抄袭并将它们混合在一起并添加到您的插件中。
  • 我肯定会从 MySQL 钩子开始,因为这样您就可以使用气流的能力来存储和检索加密的连接字符串等等。当然,您可以直接在 SQLAlchemy 中编写数据库连接和处理,但是抽象层已经存在,所以为什么不使用它。您会在其中找到已经读取到内存数据帧中的方法。父类有一个get pandas df:github.com/apache/incubator-airflow/blob/…

标签: mysql python-3.x airflow apache-airflow


【解决方案1】:

当然,只需创建一个钩子或运算符并调用 get_records() 方法:https://airflow.apache.org/docs/apache-airflow/stable/_modules/airflow/hooks/dbapi.html

【讨论】:

  • 仅供参考,链接现已修复(2019 年 6 月)
  • 链接仍然断开(2021 年 1 月)。
  • 链接仍然断开(2021 年 4 月)
【解决方案2】:

在过去的 90 分钟里,我一直在为这个问题苦苦挣扎,对于新手来说,这是一种更具说明性的方式:

from airflow.hooks.mysql_hook import MySqlHook

def fetch_records():
  request = "SELECT * FROM your_table"
  mysql_hook = MySqlHook(mysql_conn_id = 'the_connection_name_sourced_from_the_ui', schema = 'specific_db')
  connection = mysql_hook.get_conn()
  cursor = connection.cursor()
  cursor.execute(request)
  sources = cursor.fetchall()
  print(sources)

...your DAG() as dag: code

task = PythonOperator(
  task_id = 'fetch_records',
  python_callable = fetch_records
)

这会将您的数据库查询的内容返回到日志中。

我希望这对其他人有用。

【讨论】:

    猜你喜欢
    • 2018-07-16
    • 1970-01-01
    • 2016-03-14
    • 1970-01-01
    • 2021-05-01
    • 2011-02-17
    • 2019-03-31
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多