如何在Airflow中连接HDFS并执行HDFS操作?
如何在Airflow中执行HDFS操作?
1. 安装必要依赖
首先确保安装Airflow官方提供的HDFS扩展包:
pip install apache-airflow-providers-apache-hdfs
2. 配置Airflow与HDFS的连接
通过以下DAG代码在Airflow中创建HDFS连接(注意:连接成功添加到Airflow数据库后,需注释掉session.add(conn)行,避免重复创建相同连接):
# 导入依赖包 from airflow import settings from airflow.models import Connection from airflow.utils.dates import days_ago from datetime import timedelta from airflow.operators.bash import BashOperator # 定义DAG dag_execute_hdfs_commands = DAG( dag_id='connect_hdfs', schedule_interval='@once', start_date=days_ago(1), dagrun_timeout=timedelta(minutes=60), description='配置HDFS连接并执行HDFS命令', ) # 配置HDFS连接参数 conn = Connection( conn_id='webhdfs_default1', conn_type='HDFS', host='localhost', # 替换为你的HDFS主机地址 login='usr_id', # 替换为你的HDFS用户名 password='password', # 替换为你的HDFS密码 port='9000', # 替换为你的HDFS端口 ) session = settings.Session() # 将连接写入Airflow数据库,成功运行一次后必须注释此行 session.add(conn) # 后续运行前请注释掉这一行 session.close() if __name__ == '__main__': dag_execute_hdfs_commands.cli()
3. 执行HDFS操作示例
连接配置完成后,可通过BashOperator调用HDFS命令执行操作。以下是列出HDFS根目录文件的示例任务:
# 列出HDFS根目录文件的任务 start_task = BashOperator( task_id="start_task", bash_command="hdfs dfs -ls /", dag=dag_execute_hdfs_commands ) start_task
内容的提问来源于stack exchange,提问作者Swapnil
相关产品推荐
相关产品推荐

