在Apache Airflow中连接非UTF-8编码FTP服务器如何指定编码?
问题:Apache Airflow连接非UTF-8编码FTP服务器崩溃,无法修改编码配置
我需要访问一台不使用UTF-8编码的FTP服务器,使用Apache Airflow连接时会崩溃。通过Airflow底层依赖的ftplib测试发现:
- 指定
encoding='utf-8'会触发UnicodeDecodeError,报错信息如下:
ftp = FTP('myserver', user='xxxx', passwd='yyyy', encoding='utf-8') Traceback (most recent call last): File "<stdin>", line 1, in <module> File "/usr/lib/python3.11/ftplib.py", line 121, in __init__ self.connect(host) File "/usr/lib/python3.11/ftplib.py", line 162, in connect self.welcome = self.getresp() ^^^^^^^^^^^^^^ File "/usr/lib/python3.11/ftplib.py", line 244, in getresp resp = self.getmultiline() ^^^^^^^^^^^^^^^^^^^ File "/usr/lib/python3.11/ftplib.py", line 230, in getmultiline line = self.getline() ^^^^^^^^^^^^^^ File "/usr/lib/python3.11/ftplib.py", line 212, in getline line = self.file.readline(self.maxline + 1) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "<frozen codecs>", line 322, in decode UnicodeDecodeError: 'utf-8' codec can't decode byte 0xe1 in position 60: invalid continuation byte
- 指定
encoding='latin-1'则可正常连接:
ftp = FTP('myserver', user='xxxx', passwd='yyyyy', encoding='latin-1') print(ftp.welcome) 220-Microsoft FTP Service 220 FTP XXXXX, utilizado pelos usuários do orgão e Editoras.
但我在Airflow的Operator或Sensor中未找到修改编码的选项,在连接的extras中配置{"encoding":"latin-1"}也无效,请问该如何解决?
解决方案
方法1:自定义FTP Hook扩展默认实现
Airflow默认的FTPHook没有提供编码配置入口,你可以通过继承FTPHook重写连接逻辑,指定编码:
from airflow.providers.ftp.hooks.ftp import FTPHook from ftplib import FTP class Latin1FTPHook(FTPHook): def __init__(self, ftp_conn_id: str = "ftp_default"): super().__init__(ftp_conn_id=ftp_conn_id) self.encoding = "latin-1" def get_conn(self) -> FTP: conn = FTP( host=self.host, user=self.login, passwd=self.password, encoding=self.encoding ) if self.port is not None: conn.connect(port=self.port) return conn
之后在你的Operator或Sensor中使用这个自定义Hook替代默认的FTPHook即可。
方法2:适配默认Hook读取连接extras的编码配置
如果不想自定义Hook,可以修改默认FTPHook的逻辑,让它读取连接extras中的编码设置:
- 在Airflow UI的连接管理中,给目标FTP连接的extras添加
{"encoding": "latin-1"} - 在Airflow的
plugins目录下创建一个Python文件,添加以下猴子补丁代码:
from airflow.providers.ftp.hooks.ftp import FTPHook from ftplib import FTP original_get_conn = FTPHook.get_conn def patched_get_conn(self) -> FTP: # 从extras中读取编码,默认utf-8 encoding = self.extra_dejson.get("encoding", "utf-8") conn = FTP( host=self.host, user=self.login, passwd=self.password, encoding=encoding ) if self.port is not None: conn.connect(port=self.port) return conn FTPHook.get_conn = patched_get_conn
Airflow启动时会自动加载插件,之后所有使用FTPHook的地方都会读取extras中的编码配置。
方法3:直接用PythonOperator编写自定义ftplib逻辑
如果只是简单的FTP操作,也可以绕过Airflow的FTP Hook,直接在PythonOperator中编写自定义连接代码:
from airflow.operators.python import PythonOperator from ftplib import FTP def execute_ftp_operation(): # 手动指定编码连接FTP服务器 ftp = FTP('myserver', user='xxxx', passwd='yyyy', encoding='latin-1') # 执行你的FTP操作,示例:下载文件 with open('local_file.txt', 'wb') as f: ftp.retrbinary('RETR remote_file.txt', f.write) ftp.quit() ftp_task = PythonOperator( task_id="custom_ftp_task", python_callable=execute_ftp_operation )
内容的提问来源于stack exchange,提问作者Dewes
相关产品推荐
相关产品推荐

