大模型时代的“数据基座”:如何使用Spark+Hive构建亿级向量数据管道

前言

很多从大数据转向大模型开发的工程师,都会有一种割裂感:以前每天写Spark SQL、调优Hive分区,现在突然要搞Embedding、搞Prompt Engineering,感觉像是跨了个行业。但仔细看,底层逻辑其实没变。

大数据的核心是数据流转,大模型工程同样是数据流转。只不过以前输入输出是结构化表,现在输入变成了文档、日志、聊天记录,输出变成了自然语言。你在处理海量数据时的稳定性意识、容错机制、资源监控,这些在构建RAG(检索增强生成)系统时同样关键。

本文将带你走通一条完整的技术路径:从Hive离线数仓中提取业务数据 → 使用Spark进行数据清洗和向量化预处理 → 写入向量数据库供大模型应用消费


一、整体架构设计

在开始写代码之前,我们先明确一下数据流向:

Hive离线数仓(业务数据/文档语料)
    ↓
Spark ETL(数据清洗 + 文本预处理)
    ↓
Embedding生成(调用Embedding模型将文本转为向量)
    ↓
向量数据库(Milvus / LanceDB)
    ↓
RAG应用(大模型检索增强生成)

这套架构的核心思路是:把大数据工程师擅长的ETL能力,迁移到大模型时代的知识库构建场景中。以前你设计T+1的数据仓库,现在设计实时更新的知识库更新管道。这种工程决策能力,比单纯背下LangChain的API重要得多。


二、从Hive读取数据

第一步,使用Spark SQL从Hive表中读取原始数据。这里的表可以是业务工单、产品文档、客服聊天记录等任何需要作为大模型知识库的文本数据。

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, concat_ws, regexp_replace, length

# 初始化SparkSession,开启Hive支持
spark = SparkSession.builder \
    .appName("VectorDataPipeline") \
    .config("spark.sql.adaptive.enabled", "true") \
    .config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
    .enableHiveSupport() \
    .getOrCreate()

# 从Hive表读取原始数据
df_raw = spark.sql("""
    SELECT 
        doc_id,
        title,
        content,
        category,
        create_time,
        source_type
    FROM dwd.dwd_document_table
    WHERE dt = '2026-07-11'  -- 分区过滤,避免全表扫描
      AND content IS NOT NULL
      AND LENGTH(content) > 50
""")

print(f"读取数据量: {df_raw.count()}")

代码讲解

  • enableHiveSupport() 让Spark可以直连Hive Metastore,直接写SQL查询Hive表。
  • 分区过滤 dt = '2026-07-11' 是必要的,避免扫描全表导致任务崩溃。
  • LENGTH(content) > 50 过滤掉太短的无效文档。

三、数据清洗与文本预处理

做大数据时,我们最怕脏数据导致任务报错。到了大模型阶段,脏数据的定义变了——以前是空值、格式错误,现在是语义模糊、噪声干扰。

在构建RAG知识库时,如果直接把客服聊天记录或原始文档丢进模型,效果极差。因为聊天里充满了口语、缩写和情绪宣泄。这时候需要引入文本清洗策略:

from pyspark.sql.functions import udf, regexp_replace, trim
from pyspark.sql.types import StringType

# 定义文本清洗UDF
def clean_text(text):
    if text is None:
        return ""
    # 去除多余空白和换行
    text = " ".join(text.split())
    # 去除特殊字符(保留中文、英文、数字、基本标点)
    import re
    text = re.sub(r'[^\u4e00-\u9fa5a-zA-Z0-9,。!?、:;""''()\s]', '', text)
    # 截断过长的文本(防止Embedding超限)
    max_len = 2000
    if len(text) > max_len:
        text = text[:max_len]
    return text

clean_udf = udf(clean_text, StringType())

# 应用清洗
df_cleaned = df_raw \
    .withColumn("cleaned_title", clean_udf(col("title"))) \
    .withColumn("cleaned_content", clean_udf(col("content"))) \
    .withColumn("full_text", concat_ws("。", col("cleaned_title"), col("cleaned_content"))) \
    .filter(length(col("full_text")) > 20)

# 缓存清洗后的数据,方便后续复用
df_cleaned.cache()

代码讲解

  • concat_ws("。", title, content) 将标题和内容用句号拼接,形成完整的待Embedding文本。
  • 清洗策略不用一次追求完美,工程上讲究投入产出比——先做轻量级治理,后续再逐步优化。

四、批量生成Embedding向量

这是整个管道中最关键的一步:将清洗后的文本批量转化为向量。

2025-2026年的主流做法是在Spark集群中分布式调用Embedding服务,而不是在单机上循环处理。这里我展示两种方案:

方案一:使用内置Embedding服务(推荐)

如果你的公司有部署好的Embedding服务(如OpenAI Embedding API、或本地部署的BGE-M3等),可以在Executor上并发调用:

import requests
from pyspark.sql.functions import pandas_udf, struct
from pyspark.sql.types import ArrayType, FloatType
import pandas as pd
import numpy as np

# 注意:pandas_udf 需要在每个Executor上初始化HTTP连接池
EMBEDDING_URL = "http://embedding-service:8080/v1/embeddings"
BATCH_SIZE = 32  # 每批次处理32条,平衡吞吐量和延迟

@pandas_udf(ArrayType(FloatType()), functionType=pandas_udf.PandasUDFType.SCALAR)
def get_embeddings_udf(texts: pd.Series) -> pd.Series:
    """批量生成Embedding的pandas UDF"""
    # 每批处理多个文本,减少HTTP请求次数
    results = []
    session = requests.Session()
    for i in range(0, len(texts), BATCH_SIZE):
        batch = texts[i:i+BATCH_SIZE].tolist()
        try:
            resp = session.post(
                EMBEDDING_URL,
                json={"input": batch, "model": "bge-m3"},
                timeout=30
            )
            resp.raise_for_status()
            embeddings = resp.json()["data"]
            # 按原始顺序提取向量
            batch_results = [emb["embedding"] for emb in sorted(embeddings, key=lambda x: x["index"])]
            results.extend(batch_results)
        except Exception as e:
            print(f"Embedding调用失败: {e}")
            # 失败时返回空向量,后续过滤掉
            results.extend([[] for _ in batch])
    return pd.Series(results)

# 应用Embedding UDF
df_with_vectors = df_cleaned \
    .withColumn("embedding", get_embeddings_udf(col("full_text"))) \
    .filter(size(col("embedding")) > 0)  # 过滤掉失败的记录

代码讲解

  • pandas_udf 可以在Spark Executor上批量处理数据,避免Driver单点瓶颈。
  • 每批32条调用Embedding服务,在API并发限制和吞吐量之间取得平衡。
  • 生产环境中,建议在UDF内部维护连接池(如 requests.Session),减少TCP握手开销。

方案二:使用Spark原生向量化(适用于列式存储)

如果你的数据已经存储在ORC/Parquet等列式格式中,Spark的向量化查询执行可以大幅提升性能。Spark从2.0开始支持Parquet向量化读取,从2.3开始支持ORC。关键在于开启配置:

-- 开启Spark向量化读取
SET spark.sql.parquet.enableVectorizedReader=true;
SET spark.sql.orc.enableVectorizedReader=true;

-- Hive端也需要开启(如果是Hive on Spark)
SET hive.vectorized.execution.enabled=true;

当使用 df.select("embedding") 这样只读取向量列的操作时,Spark可以一次读取一批列值而不是逐行处理,CPU效率提升显著。


五、写入向量数据库

生成好的向量不能只留在Spark DataFrame里,需要导入向量数据库供RAG应用查询。这里以LanceDB为例(它原生支持Spark写入,且是云原生向量数据库的新标准):

import lancedb
from lancedb.embeddings import get_registry

# 连接LanceDB(数据存储在S3/OSS上,计算与存储分离)
db = lancedb.connect("s3://my-bucket/lancedb")

# 将Spark DataFrame转换为Pandas并写入(对于亿级数据,建议分批写入)
def write_to_lancedb(batch_df, batch_id):
    pdf = batch_df.select("doc_id", "full_text", "embedding", "category").toPandas()
    # 直接写入LanceDB表,自动创建向量索引
    table = db.create_table(
        f"doc_embeddings_{batch_id}",
        data=pdf,
        mode="overwrite",
        vector_column="embedding"
    )
    return batch_id

# 分批写入,避免OOM
df_with_vectors.foreachPartition(write_to_lancedb)

为什么选LanceDB?

  • 它基于Lance列式存储格式,原生支持向量索引 + 标量字段混合检索。
  • 索引构建可以充分利用Spark做分布式并行计算。
  • 计算与存储分离的设计,让查询节点无状态,易于水平扩展。

关于索引参数的工程调优

向量数据库的查询性能很大程度上取决于索引参数。以HNSW算法为例:

table.create_index(
    metric="cosine",
    vector_column="embedding",
    index_type="HNSW",
    index_params={
        "ef_construction": 200,  # 构建时搜索范围,越大索引质量越高但构建越慢
        "ef_search": 50,        # 查询时搜索范围,越大召回率越高但延迟越高
        "M": 16                 # 每个节点的最大连接数
    }
)

在生产环境中,建议通过压测找到平衡点。比如在我的某个客服问答项目中,将ef_search从默认的40调整到80后,召回率提升了8%,但P99延迟从45ms升到了120ms。如果你能在简历的项目描述里写出类似“通过调整向量索引参数,将召回时间控制在50ms以内”,这比罗列一堆架构名词更有说服力。


六、总结与工程建议

通过上述步骤,我们完成了一条从Hive → Spark清洗 → Embedding → 向量数据库的完整数据管道。这套架构有以下价值:

传统做法 本文方案
在Notebook中临时拼接数据 标准化ETL管道,可重复执行
单机生成向量,处理百万级数据崩溃 Spark分布式处理,轻松支撑亿级数据
向量与元数据分散存储 统一向量+标量存储,支持混合检索

三条工程建议

  1. 评估指标要量化:如果这个项目是你要写在简历上的,别只说“构建了向量数据管道”,要写“处理了X亿条历史工单数据,向量召回准确率提升X%,延迟控制在Xms以内”。

  2. 成本意识要有:Embedding服务的Token消耗是企业实打实的成本。建议对高频查询的向量结果做缓存,避免重复计算。在简历里提一句“通过缓存机制减少重复Query的Token消耗达40%”会很有分量。

  3. 监控比什么都重要:管道跑通了只是第一步。在生产环境中,你需要监控数据量波动、Embedding服务可用率、向量索引构建耗时等指标。可运维性,往往是区分初级和高级工程师的分水岭

希望这篇文章能帮助你打通从大数据到AI工程的任督二脉。如果你有具体的业务场景或技术选型问题,欢迎在评论区交流讨论!

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐