Airflow DAG调用ChromaDB插入数据时一直运行无结果求助
Airflow触发ChromaDB插入任务卡住,本地运行正常的排查与解决
问题概述
本地运行脚本可正常向ChromaDB插入数据,但通过Airflow DAG触发后,任务一直处于运行状态,无数据写入。数据库连接正常,collection.count()、collection.get()等查询类方法可正常执行,仅add和upsert写入类方法异常。
相关代码片段
Airflow调用的核心函数
def ingest_data(content): print("Received from airflow") data = content if data: doc = get_text_chunks_langchain(data) print("Sending to Embed") create_client_and_embed(doc) return "Done"
数据处理与ChromaDB操作逻辑
from langchain.text_splitter import CharacterTextSplitter import chromadb import uuid def get_text_chunks_langchain(text): text_splitter = CharacterTextSplitter(chunk_size=1000, chunk_overlap=200) docs = text_splitter.split_text(text) return docs def create_client_and_embed(docs): print(docs) print("Starting Embedding") client = chromadb.PersistentClient(path="/pathToDB") print("Printing this") # 原代码此处缺少右括号,需修正 collection = client.get_or_create_collection( name='test-db', metadata={"hnsw:space": "cosine"} ) # 任务卡在该行 collection.add(documents=docs, ids=[str(uuid.uuid4())])
Airflow DAG定义
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def callScript(**kwargs): content = kwargs.get("content") print(content) data = content return ingest_data(data) with DAG( dag_id='chroma_ingest_dag', start_date=datetime(2024, 1, 1), schedule_interval=None, catchup=False ) as dag: hello_world = PythonOperator( task_id='hello', python_callable=callScript, op_kwargs={'content': 'This is some example content to be ingested'}, provide_context=True )
排查与解决方法
- 检查文件系统权限:ChromaDB的持久化路径
/pathToDB需对Airflow worker进程所属用户(通常是airflow)开放读写权限。执行sudo chown -R airflow:airflow /pathToDB修改权限,或更换到Airflow有权限的目录。 - 同步Python环境依赖:Airflow worker的Python环境可能缺少嵌入模型相关依赖,执行
pip install sentence-transformers chromadb langchain确保依赖与本地环境一致,重点检查sentence-transformers版本。 - 处理模型自动下载阻塞:ChromaDB首次使用默认模型会自动下载,若Airflow worker无网络或下载缓慢,可提前在worker上下载模型,或显式指定本地模型路径:
from chromadb.utils import embedding_functions # 指定本地已下载的模型路径或预下载模型 embedding_func = embedding_functions.SentenceTransformerEmbeddingFunction( model_name="all-MiniLM-L6-v2", model_kwargs={'device': 'cpu'} ) collection = client.get_or_create_collection( name='test-db', metadata={"hnsw:space": "cosine"}, embedding_function=embedding_func ) - 修正代码语法错误:原代码中
get_or_create_collection调用末尾缺少右括号,会导致语法错误,需补上(如上述代码片段所示)。 - 调整任务资源与超时:在PythonOperator中增加超时设置
execution_timeout=timedelta(minutes=10),并检查Airflow worker的CPU、内存资源是否充足,必要时扩容。 - 查看任务日志:在Airflow UI中查看任务的完整日志,对比本地运行日志,排查是否有隐藏的报错信息(如权限不足、模型加载失败等)。
内容的提问来源于stack exchange,提问作者Daipayan Dey
相关产品推荐
相关产品推荐

