如何在Airflow中连接Couchbase?解决UnAmbiguousTimeoutException异常
解决Airflow PythonOperator连接Docker部署Couchbase超时问题
核心排查与解决步骤
1. 网络连通性修复
- 别用
localhost当连接地址:本地脚本能连是因为端口映射到了宿主机,但Airflow如果跑在容器里,localhost指向容器自身,不是宿主机。换成Couchbase容器的名称(前提是和Airflow容器在同一Docker网络),或者用容器IP(通过docker inspect <couchbase容器ID> | grep IPAddress获取)。 - 确认Docker网络归属:执行
docker network ls查看Couchbase所在网络,确保Airflow容器加入了这个网络,不然跨网络连不通。 - 检查端口开放:Couchbase需要开放11210(KV)、8091(管理)、8092(查询)等端口,确保Docker run命令里正确映射,且Airflow环境能访问这些端口。
2. 连接代码调整
- 加长超时时间:在连接时显式设置等待集群就绪的超时,避免因初始化慢导致超时:
from couchbase.cluster import Cluster, ClusterOptions from couchbase_core.cluster import PasswordAuthenticator auth = PasswordAuthenticator('你的用户名', '你的密码') cluster = Cluster('couchbase://couchbase容器名称', ClusterOptions(auth)) cluster.wait_until_ready(timeout=30) # 设置30秒超时,按需调整 bucket = cluster.bucket('目标桶名')
3. Airflow环境依赖检查
- 确认Couchbase SDK版本一致:Airflow Worker环境要装和本地脚本相同版本的Couchbase Python SDK,执行
pip list | grep couchbase核对版本,版本不兼容会导致各种奇怪问题。 - Docker部署的Airflow要在Worker的Dockerfile里加安装命令:
RUN pip install couchbase==<你本地用的版本号>
4. Couchbase状态验证
- 登录Couchbase管理控制台(宿主机IP:8091),确认集群正常、目标桶已创建且可用。
- 检查防火墙/安全组,确保Airflow所在IP能访问Couchbase的所有必要端口。
Airflow中Couchbase数据读写示例
假设已经从MySQL拿到DataFrame,以下是适配Airflow的读写代码(放在PythonOperator的callable函数里):
import pandas as pd from couchbase.cluster import Cluster, ClusterOptions from couchbase_core.cluster import PasswordAuthenticator def sync_mysql_to_couchbase(): # 从MySQL读取数据生成DataFrame(这里替换成你已有的读取逻辑) df = pd.read_sql("SELECT * FROM mysql_table", your_mysql_connection) # 连接Couchbase auth = PasswordAuthenticator('admin', '123456') cluster = Cluster('couchbase://couchbase-container', ClusterOptions(auth)) cluster.wait_until_ready(timeout=30) bucket = cluster.bucket('my-bucket') coll = bucket.default_collection() # 批量写入 for idx, row in df.iterrows(): doc_id = f"user_{row['id']}" coll.upsert(doc_id, row.to_dict()) # 示例读取操作 result = coll.get("user_1") print(f"读取到数据: {result.content}")
内容的提问来源于stack exchange,提问作者shubhendu
相关产品推荐
相关产品推荐

