如何通过Samba连接在Airflow中运行Python脚本?
解决方案
问题根源
SambaHook生成的路径是Samba协议格式的网络路径(如\\server\share\script.py),但Popen依赖本地操作系统的文件系统识别能力,生产环境机器未挂载该Samba共享时,系统无法直接解析这类路径,导致“找不到文件或目录”错误。
可行方案
1. 挂载Samba共享为本地路径
将远程Samba共享挂载到Airflow机器的本地文件系统,让Popen能识别本地格式的路径:
Windows环境:
用net use命令挂载,可封装为Airflow前置任务(BashOperator/PythonOperator):# 挂载共享(若未挂载) net use Z: \\your-samba-server\your-share /user:your-username your-password /persistent:yes之后脚本路径使用本地盘符格式:
Z:\path\to\transformer.py,直接传入Popen即可。Linux环境:
先安装cifs-utils,再用mount.cifs挂载:# 安装依赖(仅首次执行) sudo apt-get install cifs-utils # 挂载共享 sudo mount -t cifs //your-samba-server/your-share /mnt/samba-share -o username=your-username,password=your-password脚本路径改为本地挂载点格式:
/mnt/samba-share/path/to/transformer.py。也可写入/etc/fstab实现开机自动挂载。
2. 下载脚本到本地临时目录执行
通过SambaHook将远程脚本下载到Airflow机器的临时目录,执行本地副本后清理临时文件:
import tempfile import os from airflow.providers.samba.hooks.samba import SambaHook from subprocess import Popen, PIPE, CalledProcessError def execute_transformer_script(**context): # 初始化SambaHook samba_hook = SambaHook(connection_id="your_samba_conn_id", file_share_name="your_share_name") remote_script_path = "relative/path/to/transformer.py" # 创建临时目录 with tempfile.TemporaryDirectory() as temp_dir: local_script = os.path.join(temp_dir, os.path.basename(remote_script_path)) # 从Samba下载脚本到本地 samba_hook.retrieve_file(remote_full_path=remote_script_path, local_full_path_or_buffer=local_script) # 执行本地脚本 try: proc = Popen(["python", local_script], stdout=PIPE, stderr=PIPE, text=True) stdout, stderr = proc.communicate() proc.check_returncode() print(f"脚本执行输出: {stdout}") except CalledProcessError as e: raise RuntimeError(f"脚本执行失败: {stderr}") from e
此方案无需修改机器挂载配置,适合权限受限的生产环境,仅需注意大脚本的下载耗时。
3. 直接在Python进程中执行脚本内容
若转换脚本是可独立运行的Python模块,可读取脚本内容后直接在Airflow的Python环境中执行,避免调用外部Popen:
from airflow.providers.samba.hooks.samba import SambaHook import io def run_transformer_in_process(**context): samba_hook = SambaHook(connection_id="your_samba_conn_id", file_share_name="your_share_name") remote_script_path = "relative/path/to/transformer.py" # 读取远程脚本内容 script_buffer = io.StringIO() samba_hook.retrieve_file(remote_full_path=remote_script_path, local_full_path_or_buffer=script_buffer) script_buffer.seek(0) # 执行脚本(需确保脚本依赖已安装在Airflow环境中) exec(script_buffer.read())
注意:脚本中涉及的文件操作(如读取原始数据)需改为通过SambaHook实现,而非本地路径访问。
内容的提问来源于stack exchange,提问作者Marci
相关产品推荐
相关产品推荐

