You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 09:42:43