基于华为云FunctionGraph + OBS + GaussDB构建Serverless向量检索流水线:架构设计与代码实现
在RAG(检索增强生成)和语义搜索应用中,向量检索是一类典型的突发型工作负载。用户上传的文档需要实时进行文本嵌入、向量入库和相似度检索,这些操作的计算量随文档长度和查询复杂度波动剧烈——低峰期可能每小时只有几次上传,高峰期可能每分钟数百次并发检索。传统做法是维护一台常驻的向量数据库服务器,但资源利用率长期偏低,而在突发流量到来时又容易出现查询排队甚至超时。
华为云函数工作流FunctionGraph提供了一种更匹配这类场景的解法:按调用次数计费、毫秒级弹性伸缩、无需管理服务器。结合GaussDB的向量检索能力,可以在不维护任何常驻基础设施的前提下,构建一条完整的向量检索流水线。本文将完整实现“文档上传→OBS触发→向量化→入库→检索”的Serverless架构,包含FunctionGraph函数代码、OBS存储集成和GaussDB向量操作的完整实现。
一、整体架构
流水线由四个环节组成:用户通过API网关或OBS控制台上传文档;OBS事件触发FunctionGraph函数;函数内调用嵌入模型将文本向量化,将结果写入GaussDB的向量表;检索时通过另一个函数执行向量相似度查询。OBS作为文档存储层,GaussDB作为向量存储和检索引擎,FunctionGraph作为计算层按需运行。
二、GaussDB向量表设计与索引创建
GaussDB支持floatvector向量类型和DiskANN索引,向量检索延迟在亿级数据规模下可控制在毫秒级。首先需要在GaussDB中创建向量表和索引。以下SQL通过Python驱动执行,前提是设置GUC参数enable_vectordb=on:
import psycopg2 def init_vector_table(conn): """初始化向量表和DiskANN索引""" with conn.cursor() as cur: # 启用向量数据库功能 cur.execute("SET enable_vectordb = on") # 创建文档向量表 cur.execute(""" CREATE TABLE IF NOT EXISTS doc_vectors ( id BIGSERIAL PRIMARY KEY, doc_name VARCHAR(512) NOT NULL, chunk_index INTEGER NOT NULL, content TEXT NOT NULL, embedding FloatVector(768), created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ) """) # 创建DiskANN向量索引(余弦距离) cur.execute(""" CREATE INDEX IF NOT EXISTS idx_doc_embedding ON doc_vectors USING gsdiskann (embedding COSINE) WITH (pq_nseg = 128, pq_nclus = 16, num_parallels = 32) """) # 创建doc_name普通索引用于过滤 cur.execute(""" CREATE INDEX IF NOT EXISTS idx_doc_name ON doc_vectors (doc_name) """) conn.commit()
DiskANN索引的pq_nseg和pq_nclus参数控制乘积量化的分段数和聚类数,直接影响检索精度和速度的平衡。在768维向量的典型场景下,上述参数配置可以在召回率和延迟之间取得较好的折中。
三、FunctionGraph向量化函数实现
FunctionGraph的事件函数入口为handler(event, context),其中event包含触发器信息。OBS事件触发时,event中会包含桶名和对象名,函数需要通过OBS SDK读取文件内容,再调用嵌入模型进行向量化。以下代码使用HuggingFace的sentence-transformers作为本地嵌入模型,也可以替换为调用华为云ModelArts的远程嵌入API:
# -*- coding: utf-8 -*- import json import io import os import logging import re import psycopg2 from psycopg2 import pool from obs import ObsClient from sentence_transformers import SentenceTransformer logger = logging.getLogger() logger.setLevel(logging.INFO) # OBS客户端(全局复用) obs_client = ObsClient( access_key_id=os.getenv('OBS_AK'), secret_access_key=os.getenv('OBS_SK'), server=os.getenv('OBS_ENDPOINT') ) # 数据库连接池(全局复用) db_pool = psycopg2.pool.ThreadedConnectionPool( minconn=2, maxconn=10, host=os.getenv('DB_HOST'), port=5432, database=os.getenv('DB_NAME'), user=os.getenv('DB_USER'), password=os.getenv('DB_PASSWORD'), sslmode='require' ) # 嵌入模型(全局复用,避免每次冷启动重新加载) _model = None def get_model(): global _model if _model is None: _model = SentenceTransformer('BAAI/bge-base-zh-v1.5') return _model def handler(event, context): request_id = context.getRequestId() logger.info(f"Request {request_id} started, event: {json.dumps(event)}") # 1. 从OBS事件中提取桶名和对象名 bucket_name, object_key = _parse_obs_event(event) if not bucket_name or not object_key: return _response(400, {'error': '无法解析OBS事件'}) # 2. 从OBS读取文档内容 resp = obs_client.getObject(bucketName=bucket_name, objectKey=object_key) if resp.status >= 300: return _response(500, {'error': f'OBS读取失败: {resp.errorMessage}'}) content = resp.body.read().decode('utf-8', errors='ignore') content = re.sub(r'\s+', ' ', content).strip() if not content: return _response(400, {'error': '文档内容为空'}) # 3. 文本分块(按固定长度切分,保留语义完整性) chunks = _split_chunks(content, max_chars=500, overlap=50) logger.info(f"Document split into {len(chunks)} chunks") # 4. 向量化并写入GaussDB model = get_model() conn = db_pool.getconn() try: # 启用向量功能 with conn.cursor() as cur: cur.execute("SET enable_vectordb = on") # 批量插入 with conn.cursor() as cur: for idx, chunk in enumerate(chunks): embedding = model.encode(chunk, normalize_embeddings=True) vector_str = '[' + ','.join(str(round(float(v), 6)) for v in embedding) + ']' cur.execute( """INSERT INTO doc_vectors (doc_name, chunk_index, content, embedding) VALUES (%s, %s, %s, %s::floatvector)""", (object_key, idx, chunk, vector_str) ) conn.commit() except Exception as e: conn.rollback() logger.error(f"DB insert failed: {str(e)}") return _response(500, {'error': f'入库失败: {str(e)}'}) finally: db_pool.putconn(conn) return _response(200, { 'message': '向量化完成', 'doc_name': object_key, 'chunks': len(chunks) }) def _parse_obs_event(event): """从OBS事件中解析桶名和对象名""" try: records = event.get('Records', []) if records: s3 = records[0].get('s3', {}) bucket = s3.get('bucket', {}).get('name', '') obj = s3.get('object', {}).get('key', '') return bucket, obj except Exception as e: logger.error(f"Parse OBS event failed: {str(e)}") return None, None def _split_chunks(text, max_chars=500, overlap=50): """按语义边界切分文本块""" sentences = re.split(r'(?<=[。!?.!?])\s*', text) chunks = [] current = '' for sent in sentences: if not sent.strip(): continue if len(current) + len(sent) <= max_chars: current += sent else: if current: chunks.append(current.strip()) # 处理重叠 if overlap > 0 and len(current) > overlap: current = current[-overlap:] + sent else: current = sent if current.strip(): chunks.append(current.strip()) return chunks def _response(status_code, body): return { 'statusCode': status_code, 'headers': {'Content-Type': 'application/json'}, 'body': json.dumps(body, ensure_ascii=False) }
这段代码的关键设计点有三处。OBS事件解析:OBS触发器会将上传事件以S3兼容格式传入event,_parse_obs_event函数从中提取桶名和对象名。全局模型缓存:嵌入模型通过全局变量_model缓存,同一个函数实例的多次调用复用同一个模型,减少冷启动开销。分块策略:按句子边界切分文本,而非按固定字符数粗暴截断,保证每个chunk的语义完整性。
四、GaussDB向量检索函数
文档向量入库后,检索函数接收查询文本,将其向量化后通过GaussDB的向量相似度算子执行KNN查询:
def handler(event, context): """向量检索函数入口""" query_text = event.get('queryStringParameters', {}).get('q', '') top_k = int(event.get('queryStringParameters', {}).get('top_k', '5')) if not query_text: return _response(400, {'error': '缺少查询参数q'}) model = get_model() query_vec = model.encode(query_text, normalize_embeddings=True) vector_str = '[' + ','.join(str(round(float(v), 6)) for v in query_vec) + ']' conn = db_pool.getconn() try: with conn.cursor() as cur: cur.execute("SET enable_vectordb = on") # 余弦距离检索,按相似度降序 cur.execute(""" SELECT doc_name, chunk_index, content, 1 - (embedding <=> %s::floatvector) AS similarity FROM doc_vectors WHERE embedding <=> %s::floatvector IS NOT NULL ORDER BY embedding <=> %s::floatvector LIMIT %s """, (vector_str, vector_str, vector_str, top_k)) results = [] for row in cur.fetchall(): results.append({ 'doc_name': row[0], 'chunk_index': row[1], 'content': row[2][:200], 'similarity': round(float(row[3]), 4) }) return _response(200, {'query': query_text, 'results': results}) except Exception as e: logger.error(f"Search failed: {str(e)}") return _response(500, {'error': str(e)}) finally: db_pool.putconn(conn)
GaussDB的向量检索使用<=>操作符计算余弦距离。上述查询中使用了三次参数绑定,分别用于SELECT投影、WHERE过滤和ORDER BY排序,这是PostgreSQL协议在处理向量算子时的推荐写法。DiskANN索引会自动被查询优化器选中,无需显式指定。
五、OBS触发器配置与事件绑定
在FunctionGraph控制台为向量化函数创建OBS触发器。触发器类型选择“存储(OBS)”,选择目标桶,事件类型勾选“上传对象”事件。OBS触发器会自动将上传事件以S3兼容格式传递给函数,函数中的_parse_obs_event即可从中提取桶名和对象名。
对于检索函数,则通过API网关触发器暴露HTTPS接口。请求方法选择POST或GET,超时时间设置为30000毫秒。API网关生成的公网URL可以直接被RAG应用的前端或后端调用。
六、冷启动优化与性能考量
向量检索场景中,SentenceTransformer模型的加载是冷启动的主要耗时来源——BGE-base-zh-v1.5模型约400MB,首次加载可能需要数秒。两个优化手段:将模型文件和OBS SDK、psycopg2的依赖打包为函数层(Layer),减少代码包体积;对于检索函数这种延迟敏感的场景,考虑使用预留实例(Provisioned Concurrency)保持实例常驻,彻底消除冷启动延迟,代价是产生预留费用。
数据库连接方面,GaussDB推荐使用psycopg2连接池并启用sslmode='require'。连接池的minconn和maxconn需要根据函数并发量调整——FunctionGraph每个函数实例独立维护自己的连接池,因此maxconn不宜设置过大,否则并发实例数多时会耗尽数据库连接数。典型配置为minconn=2、maxconn=10,配合函数实例的合理内存规格。
总结
本文完整实现了基于FunctionGraph、OBS和GaussDB的Serverless向量检索流水线。核心代码覆盖了OBS事件解析、文档读取与分块、SentenceTransformer嵌入模型调用、GaussDB向量表创建与DiskANN索引、KNN检索查询的完整链路。关键工程实践包括:嵌入模型和数据库连接池的全局复用减少冷启动开销、按语义边界分块保证检索质量、DiskANN索引的乘积量化参数调优、以及OBS触发器与API网关触发器的分离设计。这一架构特别适合文档检索、知识库问答、RAG应用等场景——业务量波动大、不需要常驻基础设施、按实际调用付费,是Serverless与向量检索结合的典型落地路径。
- 点赞
- 收藏
- 关注作者
评论(0)