基于Doris与AI构建非结构化数据分析系统:从向量化到智能洞察
1. 项目概述:当非结构化数据遇上AI分析
最近几年,数据领域一个明显的趋势是,我们手里的“原材料”越来越“杂”了。过去,我们处理的大多是规规矩矩的订单表、用户日志,这些结构化数据就像超市里包装好的蔬菜,清洗、切配都方便。但现在,大量的图片、PDF文档、音频、视频、网页内容涌了进来,这些非结构化数据就像刚从地里挖出来的、带着泥的土豆,形态各异,处理起来费时费力。很多团队都卡在了第一步:怎么把这些“泥土豆”洗干净、切好,变成能下锅分析的“数据食材”?
这正是“从零搭建非结构化数据智能分析洞察系统”这个项目要解决的核心问题。它不是一个空中楼阁的概念,而是一个结合了现代数据栈(Modern Data Stack)中流批处理、向量化技术与AI能力的实战工程。简单来说,就是构建一个管道,把散乱的非结构化内容(比如公司内部堆积如山的合同PDF、产品讨论会的录音、社交媒体上的图片)自动转化成结构化的、可被数据库高效查询和分析的信息,最终通过BI工具或应用,让业务人员能像查销售报表一样,去“查询”和“洞察”这些非结构化内容里的价值。
在这个方案里,Doris(或它的商业发行版SelectDB)扮演了至关重要的“中枢大脑”角色。它不再仅仅是一个传统的OLAP数据库,而是成为了一个能够统一承载原始文本、处理后的结构化特征、以及AI生成的向量化表征的“统一分析平台”。你可以把它理解为一个超级厨房,既能存放原始食材(原始文件路径、元数据),也能存放切好的配菜(解析出的文本、关键信息),还能存放用特殊调料腌制过的半成品(向量Embedding),并且能根据你的需求,快速组合出不同的“菜品”(即席查询与分析)。
这个系统的价值在于,它极大地降低了非结构化数据的使用门槛。市场团队可以快速分析竞品发布会视频的文本,提炼核心卖点;法务团队可以批量审查合同,自动识别关键条款和风险点;客服团队可以从海量录音中,定位用户抱怨的高频问题。整个过程,从数据接入、AI处理到分析洞察,形成了一个自动化闭环。
2. 系统核心架构与设计思路拆解
搭建这样一个系统,关键在于设计一个高内聚、低耦合的流水线,确保数据流顺畅、处理环节可扩展、最终查询高效。一个经过实战检验的典型架构可以分为四层:数据接入与预处理层、AI能力处理层、数据存储与分析层、以及应用与洞察层。
2.1 整体架构设计:流水线思维
我们的核心设计思路是“流水线化”和“统一入口”。整个系统像一条智能装配线:
- 原始数据投递:各种来源的非结构化文件(对象存储中的图片、消息队列里的文档链接、直接上传的文件)被统一收集到一个接入点。这里,Apache Kafka或AWS S3等对象存储是常见选择,它们负责数据的缓冲和暂存。
- AI处理车间:这是系统的“智能核心”。我们使用Flink或Spark这类流批一体处理框架,构建处理作业。这个作业会从接入点消费数据,然后调用一系列AI服务。例如,先用OCR服务处理图片,再用NLP服务对OCR文本进行实体识别和情感分析,最后调用文本嵌入模型(如
text-embedding-3-small)将文本转化为向量。这个过程可能是多步骤的、有分支的。 - 统一存储与查询枢纽:处理后的数据需要落地。这里就是Doris/SelectDB的主场。我们会将三类数据写入Doris:
- 结构化元数据:文件ID、来源、处理时间、从内容中提取的关键属性(如合同金额、签署方、文档类型)。
- 原始/处理后文本:OCR识别出的完整文本或经过清洗的文本。
- 向量数据:文本通过Embedding模型生成的、用于相似性搜索的高维向量。 将这三者放在同一张表或通过主键关联的多张表中,是实现“统一分析”的物理基础。
- 应用与洞察终端:分析师或应用程序通过标准SQL或BI工具(如FineBI、Metabase)连接Doris。他们可以:
- 用SQL做传统的条件过滤、聚合分析(“统计上个月所有合同中涉及‘保密协议’条款的数量”)。
- 利用Doris内置的向量搜索函数(如
dot_product,cosine_distance),通过自然语言进行语义搜索(“查找与‘员工股权激励方案’最相似的文档”)。 - 将两者结合,实现混合查询(“在2023年Q3的销售报告中,找出与‘市场增长乏力’描述相似的段落”)。
设计思路核心:这个架构的优势在于,它将复杂的AI处理流程封装在了数据管道中,对上游数据源和下游应用透明。下游用户无需关心文本是如何从图片里来的,也无需手动调用AI模型,他们面对的是一个已经“增强”了的、支持向量搜索的数据库表。Doris的统一性避免了数据在多个系统(如传统数据库+向量数据库)间搬运带来的延迟和一致性难题。
2.2 为什么选择Doris/SelectDB作为核心?
面对非结构化数据分析的场景,市面上有专门的向量数据库(如Milvus, Pinecone),也有传统的数仓。为什么我们倾向于选择Doris/SelectDB呢?这源于几个关键的技术决策点:
统一分析平台,避免数据碎片化:这是最重要的原因。如果使用“传统数仓+独立向量库”的架构,业务查询一个简单的“混合条件过滤+语义相似度排序”需求,就需要跨系统联合查询,复杂度高、性能差、一致性难保证。Doris通过其
ARRAY<VECTOR>数据类型和向量函数,将向量搜索变成了一个SQL函数调用,实现了在单一数据库内完成所有分析。数据只需存储一份,维护成本大大降低。卓越的实时分析与高并发性能:Doris的MPP架构和列式存储引擎,使其在复杂聚合查询、多表关联上具有传统向量数据库难以比拟的优势。当你的洞察需求不仅仅是“找相似”,还包括“按部门、时间聚合分析相似文档的分布”时,Doris的性能优势就凸显出来了。SelectDB在此基础上,进一步优化了云原生部署和向量检索性能。
成熟的生态与运维体系:Doris与大数据生态(Flink, Spark, Kafka)的集成非常成熟,数据导入(
routine load,spark connector)方式丰富。其监控、备份、扩容等运维操作,对于数据团队来说,比运维一个较新的专用向量数据库学习成本更低,体系更完善。成本与效率的平衡:对于很多企业,非结构化数据分析是新兴场景,但并非唯一场景。单独引入和维护一套向量数据库,会增加基础设施和团队的复杂度。利用现有Doris集群的能力进行扩展,是一种更务实、性价比更高的选择,尤其当数据量在千万到百亿级别,且需要复杂分析时。
当然,这个选择并非绝对。如果业务场景极端聚焦于海量向量的、低延迟的、纯相似性检索(如大规模推荐系统召回),专用向量数据库可能更合适。但对于大多数寻求“非结构化数据智能洞察”的企业级应用,Doris/SelectDB提供的“一站式”解决方案,在功能、性能和复杂度上取得了更好的平衡。
3. 核心模块实操:从数据接入到向量入库
理论讲完,我们进入实战环节。假设我们的场景是分析企业内部大量的产品评审会议纪要(PDF和音频转录文本),目标是构建一个可以按主题、情感、部门进行检索和统计的系统。
3.1 环境准备与工具选型
首先,我们需要搭建好基础设施。这里给出一个基于开源组件的推荐方案:
- 数据处理引擎:Apache Flink。我们选择Flink而非Spark Streaming,主要是看中其真正的流处理能力和更低的端到端延迟,这对于希望近实时感知非结构化数据内容的场景(如舆情监控)很重要。使用Flink 1.17+版本。
- AI模型服务化:Model Server。不建议在Flink作业内直接加载大模型,这会导致资源管理混乱和性能瓶颈。应将AI模型(OCR、NLP、Embedding)部署为独立的服务。推荐使用Ray Serve或Triton Inference Server,它们专为高性能模型推理设计,支持动态批处理、多模型版本管理,能显著提高GPU利用率和吞吐量。
- 向量生成模型:文本嵌入模型。这是向量质量的关键。对于中文场景,
BAAI/bge-large-zh或moka-ai/m3e-base是经过广泛验证的优秀开源模型。对于多语言或对效果有极致要求,可以考虑OpenAI的text-embedding-3-small(需API调用)。我们选择m3e-base,因其在中文语义相似度任务上表现好,且模型尺寸适中,便于部署。 - 核心分析平台:Apache Doris 2.0+。2.0版本对向量功能进行了大幅增强。我们需要一个至少3个BE(后端节点)的集群。如果追求更便捷的云上管理和更优的向量性能,可以直接使用SelectDB Cloud。
实操心得一:模型服务与资源隔离将AI模型推理独立部署,是一个关键架构决策。我曾尝试在Flink作业中内嵌PyTorch模型,结果发现:1)作业重启慢,因为要重新加载模型;2)GPU内存管理复杂,容易导致容器崩溃;3)无法利用动态批处理来优化吞吐。使用Ray Serve后,模型服务自成一体,可以独立扩缩容,Flink作业只需通过HTTP/gRPC调用,解耦后系统稳定性和可维护性大幅提升。
3.2 构建Flink AI处理流水线
这是整个系统的“发动机”。Flink作业负责编排整个处理逻辑。我们使用Java/Scala编写,但核心处理可能调用Python服务。
// 简化示例:Flink Job 主逻辑骨架 DataStream<String> sourceStream = env.addSource(new KafkaSource<>(...)); // 从Kafka读取文件URL或元数据 // 处理链 SingleOutputStreamOperator<EnrichedDocument> processedStream = sourceStream .map(new FetchRawContentMapper()) // 1. 根据URL获取原始文件内容(PDF/音频二进制流) .flatMap(new DocumentSplitter()) // 2. 拆分文档(如按PDF页面、按音频分钟分段) .process(new AIServiceProcessFunction()) // 3. 核心AI处理:调用OCR/ASR -> NLP -> Embedding服务 .name("AI-Enrichment-Pipeline"); // 写入Doris processedStream.addSink( DorisSink.sink( DorisExecutionOptions.builder().setBatchSize(1000).setMaxRetries(3).build(), DorisOptions.builder() .setFenodes("FE_HOST:8030") .setTableIdentifier("db.doc_ai_insights") .setUsername("user").setPassword("passwd").build(), new DocumentToRowDataSerializer() // 自定义序列化器,将对象转为Doris Row ) );关键环节详解:AIServiceProcessFunction
这个ProcessFunction是流水线的核心,它需要完成以下步骤,并做好错误处理和状态管理:
- 内容提取:判断文档类型,如果是PDF,调用部署在Ray Serve上的PaddleOCR服务,获取文本和位置信息;如果是音频,调用Whisper语音转文本服务。
- 文本增强与清洗:对提取的文本进行基础清洗(去除无意义字符、换行符规范化)。然后,调用NLP服务进行:
- 命名实体识别(NER):提取人名、公司名、产品名、日期、金额等。
- 情感分析:判断该段文本的情感倾向(正面、负面、中性)。
- 关键短语/主题提取:使用
TextRank或KeyBERT算法提取核心关键词。
- 向量化:将清洗后的完整文本或关键句子,发送给
m3e-base模型服务,获取768维的文本向量。 - 结果组装:将原始文本、提取的实体(以JSON格式存储)、情感标签、关键词列表、向量数组,以及原始文件的元数据,组装成一个
EnrichedDocument对象。
实操心得二:处理粒度与性能权衡对于长文档(如几十页的PDF),是整文档生成一个向量,还是分页/分段生成多个向量?这需要权衡。整文档向量丢失了细节,但存储和查询成本低。分段向量更精细,支持更准确的段落级搜索,但数据量会膨胀。我们的经验是:两级向量化。既为整个文档生成一个“概要向量”,也为每个逻辑段落(如PDF的每一节)生成“细节向量”。在Doris中可以用两个
ARRAY<VECTOR>字段存储。查询时,先通过概要向量快速筛选相关文档,再用细节向量进行精排,兼顾了召回速度和精度。
3.3 Doris表设计与数据导入
处理好的数据要高效地存入Doris。表结构设计直接影响查询的灵活性和性能。
CREATE DATABASE IF NOT EXISTS ai_insights; USE ai_insights; CREATE TABLE doc_ai_insights ( `doc_id` VARCHAR(255) NOT NULL, `source_type` VARCHAR(50) COMMENT 'pdf, audio, image', `file_path` VARCHAR(1000), `upload_time` DATETIME DEFAULT CURRENT_TIMESTAMP, `raw_text` STRING COMMENT '原始提取文本', `clean_text` STRING COMMENT '清洗后文本', `entities` JSON COMMENT '提取的实体,如{"person":[...], "company":[...]}', `sentiment` TINYINT COMMENT '情感分值,-1负,0中,1正', `keywords` ARRAY<VARCHAR(200)> COMMENT '关键词列表', `summary_vector` ARRAY<FLOAT> COMMENT '文档概要向量', `segment_vectors` ARRAY<ARRAY<FLOAT>> COMMENT '段落向量数组', `segment_texts` ARRAY<STRING> COMMENT '对应段落文本' ) ENGINE=OLAP DUPLICATE KEY(`doc_id`, `source_type`, `upload_time`) DISTRIBUTED BY HASH(`doc_id`) BUCKETS 10 PROPERTIES ( "replication_num" = "3", "storage_format" = "V2" );设计要点解析:
- 数据类型选择:
ARRAY<FLOAT>用于存储向量。JSON类型灵活存储非标准化的实体信息。ARRAY<ARRAY<FLOAT>>用于存储多个段落向量,这是一个数组嵌套数组的结构。 - 分区与分桶:我们以
upload_time作为分区字段(实际创建时需使用PARTITION BY RANGE),便于按时间范围管理数据。DUPLICATE KEY指定了前缀索引列,查询经常按doc_id或source_type过滤,将其放在前面能加速查询。 - 向量索引:对于向量搜索,仅靠
ARRAY<FLOAT>类型是不够的。Doris 2.1+支持倒排索引(INVERTED INDEX)和向量索引(VECTOR INDEX)。我们需要对summary_vector和展开后的segment_vectors建立向量索引以加速相似性搜索。-- 为概要向量列添加向量索引(假设使用HNSW算法) ALTER TABLE doc_ai_insights ADD INDEX vec_idx_summary(summary_vector) USING VECTOR; -- 注意:对嵌套数组建立索引可能需要将数据扁平化到另一张表,这里简化示意
数据导入使用Flink-Doris-Connector,如上节代码所示,它能保证Exactly-Once语义,确保数据不丢不重。对于存量数据的批量导入,可以使用Spark-Doris-Connector或Broker Load。
4. 智能查询、分析与应用实战
数据就位后,最激动人心的部分来了:如何查询和挖掘其中的价值?Doris的SQL能力让我们可以玩出很多花样。
4.1 基础语义搜索:用自然语言找文档
这是最直接的应用。假设你想找和“三季度销售额下滑原因分析”相关的会议纪要。
-- 首先,将查询语句也转化为向量。这里需要在应用层调用相同的Embedding模型,得到查询向量。 -- 假设我们已得到查询向量 `query_vec` (一个FLOAT数组) SET query_vec = '[0.123, -0.456, ..., 0.789]'; -- 使用余弦相似度进行搜索,并与其他条件结合 SELECT doc_id, file_path, clean_text AS content_preview, cosine_distance(summary_vector, ${query_vec}) AS distance, -- 距离越小越相似 sentiment, keywords FROM doc_ai_insights WHERE source_type = 'audio' -- 可以结合业务条件过滤 AND upload_time >= '2024-01-01' ORDER BY distance ASC -- 按相似度升序排列 LIMIT 10;4.2 混合检索:语义+条件过滤
单纯的向量搜索可能返回一些时间久远或不相关的文档。结合业务元数据进行过滤,能让结果更精准。
-- 查找市场部上传的、情感倾向为负面的、且与“用户流失”相关的文档 SELECT doc_id, -- 使用dot_product计算相似度(需向量已归一化),值越大越相似 dot_product(summary_vector, ${query_vec_loss}) AS similarity, entities['department'] AS dept, sentiment, SUBSTRING(clean_text, 1, 200) AS snippet FROM doc_ai_insights WHERE JSON_EXTRACT(entities, '$.department') = 'marketing' AND sentiment = -1 AND array_contains(keywords, '用户流失') -- 关键词过滤 ORDER BY similarity DESC LIMIT 20;4.3 聚合分析与洞察报表
将AI提取的信息进行聚合,可以生成强大的分析报表。
-- 按部门统计负面情感文档的主题分布 SELECT JSON_EXTRACT(entities, '$.department') AS department, keyword, COUNT(*) AS negative_count, COUNT(DISTINCT doc_id) AS unique_docs FROM doc_ai_insights, UNNEST(keywords) AS keyword_tbl(keyword) -- 展开关键词数组 WHERE sentiment = -1 AND upload_time BETWEEN '2024-04-01' AND '2024-04-30' GROUP BY department, keyword HAVING negative_count > 5 ORDER BY department, negative_count DESC; -- 分析某个产品名称在不同会议中被提及的情感趋势 SELECT DATE_TRUNC('week', upload_time) AS week, AVG(sentiment) AS avg_sentiment_score, COUNT(*) AS mention_times FROM doc_ai_insights WHERE array_contains(keywords, '产品A') -- 或使用JSON_EXTRACT在entities中查找 GROUP BY week ORDER BY week;这些SQL查询可以直接在Doris的Web UI中执行,也可以轻松集成到FineBI、Metabase等BI工具中,制作成可视化的仪表盘,让业务人员随时查看非结构化数据中的洞察。
4.4 性能优化关键点
当数据量增长后,查询性能至关重要。以下几点是优化关键:
- 向量索引调优:为
summary_vector创建VECTOR INDEX时,需要选择合适算法(如HNSW)和参数(metric_type=cosine,m=16,ef_construction=200)。这些参数影响索引构建速度、内存占用和检索精度,需要在你的数据集上进行测试调优。 - 分区与分桶策略:根据查询模式调整分区键。如果经常按时间范围查询,
upload_time作为分区键非常有效。分桶数量建议是BE节点数的整数倍,且单个桶数据量在100MB-1GB为宜。 - 物化视图预聚合:对于上方的聚合分析查询,如果数据量大且查询频繁,可以为
(department, sentiment, week)创建带聚合的物化视图,将计算提前,极大加速查询。 - 查询模式:避免对
segment_vectors这种嵌套数组列直接进行全表扫描的向量计算。应先通过summary_vector或条件过滤缩小数据集,再对子集进行精细的向量计算。
5. 常见问题、排查技巧与演进思考
在实际搭建和运营这套系统的过程中,你一定会遇到各种坑。这里分享一些典型的“踩坑”经验和排查思路。
5.1 典型问题与解决方案速查表
| 问题现象 | 可能原因 | 排查步骤与解决方案 |
|---|---|---|
| Flink作业消费Kafka延迟高 | 1. AI模型服务调用慢。 2. 并行度设置不合理。 3. Checkpoint时间过长。 | 1. 监控Ray Serve服务P99延迟,优化模型批处理大小。 2. 增加Flink作业并行度,特别是 AIServiceProcessFunction之后的算子。3. 调大Checkpoint间隔,或使用增量Checkpoint。 |
| 向量相似度搜索结果不相关 | 1. Embedding模型与领域不匹配。 2. 文本预处理不当(噪音多)。 3. 查询向量生成方式不一致。 | 1. 在业务数据上微调Embedding模型,或更换更适配的模型(如从m3e-base换为bge-large-zh)。2. 加强文本清洗:去停用词、特殊字符,保留关键实体。 3. 确保查询文本与入库文本使用完全相同的预处理流程和模型。 |
| Doris向量搜索速度慢 | 1. 未建立向量索引。 2. 查询时未有效利用索引。 3. 内存不足。 | 1. 确认ALTER TABLE ... ADD INDEX已成功执行,使用SHOW INDEX FROM table查看。2. 检查SQL,确保向量距离计算函数(如 cosine_distance)的参数之一是索引列。3. 监控BE节点内存,调整 vector_index_ram_limit参数。 |
| 写入Doris吞吐量低 | 1. Sink批次设置过小。 2. Doris BE节点写入压力大。 3. 网络延迟。 | 1. 调大Flink Doris Connector的batch.size和batch.interval。2. 观察Doris BE的 write_bytes_rate等指标,考虑增加BE节点或分桶。3. 确保Flink JobManager/TaskManager与Doris集群网络通畅。 |
| 提取的实体或关键词质量差 | 1. NLP模型在垂直领域表现不佳。 2. 原始文本质量差(OCR错误多)。 | 1. 使用领域词典微调NER模型,或采用规则+模型结合的方式后处理。 2. 优化OCR前处理(图像增强)和后处理(词典校正)。 |
5.2 系统演进与扩展思考
当这个基础系统跑起来后,可以考虑以下几个方向的深化:
- 多模态融合:当前以文本为主。可以扩展支持图像向量(使用CLIP模型)和音频向量,在Doris中实现“以图搜图”、“根据描述找图片”或“找包含特定声音的片段”。
- 实时流式洞察:将处理链路延迟优化到秒级,实现对直播字幕、即时通讯消息的实时情感分析和热点发现,用于舆情监控或实时客服质检。
- Agentic AI集成:将Doris作为智能体的“记忆中枢”或“知识库”。让大语言模型(LLM)通过插件或函数调用,执行上述的混合检索SQL,获取最相关的背景信息,再生成回答或报告,实现真正的“数据对话”。
- 成本优化:非结构化数据处理和向量生成是计算密集型的。可以考虑:a) 对“热数据”进行全流程处理,对“冷数据”只做基础提取和存储,需要时再异步生成向量;b) 使用量化技术压缩向量维度,在精度损失可接受范围内大幅节省存储和计算资源。
从我个人的实践经验来看,这套系统的最大挑战往往不在技术本身,而在于领域知识的注入。一个通用的Embedding模型在医疗合同和科技新闻上的表现天差地别。成功的秘诀在于,深入业务,用高质量的、经过业务标注的数据去微调每一个AI环节(特别是Embedding和NER模型),让整个系统真正“懂”你的数据。这就像给一个聪明的助手提供了专业的行业词典,它才能给出真正有价值的洞察。