This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12488-0e1a8b6f848e9cdf1dae4881d70d61103fd13cff in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 4fd80db52f401fdd47b826b1ec792d1a6cea1ac4 Author: inchang-ing <[email protected]> AuthorDate: Sun Sep 27 11:07:56 2026 +0000 [Docs] Add RAG ready data processing guide (#12488) Signed-off-by: inchang-ing <[email protected]> Co-authored-by: inchang-ing <[email protected]> --- docs/en/introduction/rag-data-processing.md | 181 ++++++++++++++++++++++++++++ docs/sidebars.js | 1 + docs/zh/introduction/rag-data-processing.md | 181 ++++++++++++++++++++++++++++ 3 files changed, 363 insertions(+) diff --git a/docs/en/introduction/rag-data-processing.md b/docs/en/introduction/rag-data-processing.md new file mode 100644 index 0000000000..dd40132939 --- /dev/null +++ b/docs/en/introduction/rag-data-processing.md @@ -0,0 +1,181 @@ +# RAG-Ready Data Processing + +> How to use SeaTunnel to build Retrieval-Augmented Generation (RAG) data pipelines: parse documents, split text into chunks, generate embeddings, and load vector databases — in batch or streaming. + +## Why a data pipeline matters for RAG + +Retrieval-Augmented Generation (RAG) is only as good as the data behind it. The path from raw knowledge — files, wikis, databases — to a vector database involves several stages: parsing documents into text, splitting text into chunks, optionally enriching or cleaning it with an LLM, converting chunks into embeddings, and loading them into a vector store that stays fresh over time. + +SeaTunnel covers all of these stages with its existing plugin model, so the same engine that moves your business data can also build and maintain your knowledge base: the same configuration style, the same parallelism and fault-tolerance behavior, and the same batch/streaming execution modes. + +## Building blocks + +| Stage | SeaTunnel plugin | Documentation | +| --- | --- | --- | +| Document parsing | LocalFile source with `file_format_type = markdown` or `pdf` | [LocalFile](../connectors/source/LocalFile.md) | +| Database ingestion | JDBC sources such as MySQL | [MySQL](../connectors/source/Mysql.md) | +| Change data capture | MySQL-CDC source | [MySQL-CDC](../connectors/source/MySQL-CDC.md) | +| Text splitting | TextChunk transform | [TextChunk](../transforms/text-chunk.md) | +| Simple field splitting | Split transform | [Split](../transforms/split.md) | +| Metadata projection | Metadata transform | [Metadata](../transforms/metadata.md) | +| LLM processing | LLM transform | [LLM](../transforms/llm.md) | +| Vectorization | Embedding transform | [Embedding](../transforms/embedding.md) | +| Vector storage | Milvus / Qdrant sinks | [Milvus](../connectors/sink/Milvus.md), [Qdrant](../connectors/sink/Qdrant.md) | + +## Pipeline at a glance + +``` + documents / tables + │ + ▼ + parse (LocalFile: markdown / pdf) ── or ── JDBC / CDC source + │ │ + ▼ ▼ + chunk (TextChunk) ────────► enrich (LLM, Metadata) ───► embed (Embedding) + │ + ▼ + vector store (Milvus / Qdrant) +``` + +## Example 1: build a document knowledge base in Milvus + +A batch job that parses a directory of Markdown documents, re-chunks the text with overlap, embeds each chunk, and writes everything to Milvus: + +```hocon +env { + parallelism = 2 + job.mode = "BATCH" +} + +source { + LocalFile { + path = "/data/knowledge-base" + file_format_type = "markdown" + markdown_rag_metadata_enabled = true + } +} + +transform { + TextChunk { + text_field = "text" + output_field = "chunk" + chunk_index_field = "chunk_seq" + chunk_size = 800 + overlap_size = 100 + } + Embedding { + model_provider = OPENAI + api_key = "sk-xxxx" + model = "text-embedding-3-small" + dimension = 1024 + single_vectorized_input_number = 10 + model_retry_max_attempts = 5 + vectorization_fields { + chunk = "chunk_vector" + } + } +} + +sink { + Milvus { + url = "http://127.0.0.1:19530" + token = "root:Milvus" + collection = "knowledge_base" + enable_upsert = false + batch_size = 500 + } +} +``` + +What happens at each stage: + +1. **Parse.** The LocalFile source parses each Markdown file into element rows (headings, paragraphs, list items, code blocks, ...). Each row carries fields such as `text`, `element_id`, `element_type`, `heading_level` and `position_index`. With `markdown_rag_metadata_enabled = true` the source additionally appends RAG metadata columns: `source_uri`, `document_id`, `chunk_id`, `chunk_index` and `content_hash` — stable identifiers derived from the file path and content. +2. **Chunk.** The TextChunk transform splits each element's `text` into chunks of at most 800 characters, keeping about 100 characters of overlap between adjacent chunks. Every output row keeps all original fields and adds the chunk text (`chunk`) and its position within the source row (`chunk_seq`). +3. **Embed.** The Embedding transform calls the embedding model and writes the vector into a new field `chunk_vector` with dimension 1024. +4. **Load.** The Milvus sink creates the target collection from the upstream schema if it does not exist and writes in batches of 500 records. + +> Note: `chunk_index` produced by the Markdown parser is 1-based, while the index emitted by TextChunk (`chunk_index_field`) is 0-based. This example renames TextChunk's field to `chunk_seq` so the two never collide. + +To ingest PDF files instead, use `file_format_type = "pdf"` together with `pdf_rag_metadata_enabled = true`. + +## Example 2: vectorize database rows into Qdrant + +The same pipeline pattern works for structured data — here, product descriptions are embedded and written to Qdrant: + +```hocon +env { + parallelism = 4 + job.mode = "BATCH" +} + +source { + MySQL { + url = "jdbc:mysql://127.0.0.1:3306/shop" + driver = "com.mysql.cj.jdbc.Driver" + username = "root" + password = "******" + table_path = "shop.products" + } +} + +transform { + Embedding { + model_provider = OPENAI + api_key = "sk-xxxx" + model = "text-embedding-3-small" + dimension = 1024 + vectorization_fields { + description = "description_vector" + } + } +} + +sink { + Qdrant { + collection_name = "product_vectors" + host = "127.0.0.1" + port = 6334 + } +} +``` + +## Keeping vectors fresh + +Knowledge bases are living systems: documents are edited and database rows change. SeaTunnel supports several refresh strategies: + +- **Full rebuild.** Re-run the job with `data_save_mode = DROP_DATA` on the Milvus sink to drop and rebuild the collection from scratch. Simple and always consistent, at the cost of a full re-embedding run. +- **Idempotent updates.** If the upstream schema carries a primary key, set `enable_upsert = true` on the Milvus sink so re-running the job updates existing records instead of duplicating them. +- **Change data capture.** Use the MySQL-CDC source to stream row changes and apply them to the vector store continuously. Pair CDC with an upsert-capable sink so updated rows are re-embedded and overwritten in place. +- **Chunk-level lifecycle.** Be aware that SeaTunnel does not automatically delete stale chunks when a source document changes (no tombstones). If you re-chunk documents and the number of chunks shrinks, old chunks remain. Either rebuild periodically, or design stable chunk identifiers and handle deletions in a follow-up step of your own pipeline. + +## Best practices + +- **Chunk size.** `chunk_size` is counted in Unicode code points and defaults to 1000. Keep `overlap_size` well below `chunk_size` — the chunk count grows roughly as `chunk_size / (chunk_size - overlap_size)`. Use `separators` to avoid cutting mid-sentence. +- **Embedding reliability.** Remote embedding APIs are rate-limited and occasionally fail. Set `model_retry_max_attempts` greater than 1 so retryable failures (rate limiting, timeouts) are retried with backoff instead of failing the job. Use `single_vectorized_input_number` to batch multiple inputs into one request when your provider allows it. +- **Precision.** Embedding vectors are produced in `float32` only; other precisions are converted automatically. Make sure the vector dimension configured on the Embedding transform matches the vector field expected by the target collection. +- **Primary keys.** TextChunk requires that `output_field` / `chunk_index_field` do not collide with columns participating in a primary or unique key. When upserting into Milvus, remember that upsert requires a primary key in the collection schema. +- **Metadata filtering.** With `markdown_rag_metadata_enabled = true`, each row carries `document_id`, `chunk_id` and `content_hash`. Use the Metadata transform to project the logical knowledge-sync fields (for example `SourceUri`, `DocumentId`, `ChunkHash`) into application-specific columns for filtering at retrieval time: + + ```hocon + Metadata { + metadata_fields = { + SourceUri = "ks_source_uri" + DocumentId = "ks_document_id" + ChunkHash = "chunk_hash" + } + } + ``` + +## FAQ + +### Can SeaTunnel delete stale chunks when a document changes? + +Not automatically. SeaTunnel moves and transforms data; chunk deletion policies belong to the consuming pipeline. See [Keeping vectors fresh](#keeping-vectors-fresh) for the supported strategies. + +### Which embedding model providers are supported? + +Common providers such as `OPENAI`, `AMAZON`, `DOUBAO` and `QIANFAN` are supported, plus a `CUSTOM` provider where you define request headers, body and response parsing yourself. See the [Embedding](../transforms/embedding.md) documentation for the full list and options. + +### Do I have to re-embed everything when one document changes? + +No — with `markdown_rag_metadata_enabled = true` each chunk carries a `content_hash` and stable identifiers, so downstream systems can detect which chunks changed. Note that if a transform between source and sink changes the text or expands rows (like TextChunk), the final chunk identifiers and hashes should be recomputed after that transform. diff --git a/docs/sidebars.js b/docs/sidebars.js index 20b6f38a1a..05cfd3e4c7 100644 --- a/docs/sidebars.js +++ b/docs/sidebars.js @@ -26,6 +26,7 @@ const sidebars = { "items": [ "introduction/about", "introduction/how-it-works", + "introduction/rag-data-processing", { "type": "category", "label": "Concepts", diff --git a/docs/zh/introduction/rag-data-processing.md b/docs/zh/introduction/rag-data-processing.md new file mode 100644 index 0000000000..00307c0088 --- /dev/null +++ b/docs/zh/introduction/rag-data-processing.md @@ -0,0 +1,181 @@ +# RAG Ready 数据处理 + +> 介绍如何使用 SeaTunnel 构建 RAG(Retrieval-Augmented Generation,检索增强生成)数据管道:解析文档、切分文本、生成向量,并写入向量数据库 —— 支持批式与流式。 + +## 为什么 RAG 需要一条数据管道 + +RAG 的效果取决于它背后的数据。从原始知识(文件、Wiki、数据库)到向量数据库,通常要经过多个阶段:把文档解析成文本、把文本切分成块、(可选)用 LLM 做清洗或增强、把文本块转换为向量,再把向量写入一个能持续保持新鲜的向量库。 + +SeaTunnel 用它现有的插件体系覆盖了上述全部阶段 —— 那套搬运业务数据的引擎,同样可以用来构建和维护你的知识库:一样的配置风格、一样的并行与容错语义、一样的批式/流式执行模式。 + +## 构建模块 + +| 阶段 | SeaTunnel 插件 | 文档 | +| --- | --- | --- | +| 文档解析 | LocalFile source,`file_format_type = markdown` 或 `pdf` | [LocalFile](../connectors/source/LocalFile.md) | +| 数据库接入 | MySQL 等 JDBC source | [MySQL](../connectors/source/Mysql.md) | +| 变更数据捕获 | MySQL-CDC source | [MySQL-CDC](../connectors/source/MySQL-CDC.md) | +| 文本切分 | TextChunk transform | [TextChunk](../transforms/text-chunk.md) | +| 简单字段拆分 | Split transform | [Split](../transforms/split.md) | +| 元数据投影 | Metadata transform | [Metadata](../transforms/metadata.md) | +| LLM 处理 | LLM transform | [LLM](../transforms/llm.md) | +| 向量化 | Embedding transform | [Embedding](../transforms/embedding.md) | +| 向量存储 | Milvus / Qdrant sink | [Milvus](../connectors/sink/Milvus.md),[Qdrant](../connectors/sink/Qdrant.md) | + +## 管道一览 + +``` + 文档 / 数据表 + │ + ▼ + 解析(LocalFile:markdown / pdf)── 或 ── JDBC / CDC source + │ │ + ▼ ▼ + 切分(TextChunk)────► 增强(LLM、Metadata)──► 向量化(Embedding) + │ + ▼ + 向量库(Milvus / Qdrant) +``` + +## 示例一:在 Milvus 中构建文档知识库 + +一个批式作业:解析 Markdown 文档目录 → 带重叠地重新切分文本 → 对每个文本块做向量化 → 全部写入 Milvus: + +```hocon +env { + parallelism = 2 + job.mode = "BATCH" +} + +source { + LocalFile { + path = "/data/knowledge-base" + file_format_type = "markdown" + markdown_rag_metadata_enabled = true + } +} + +transform { + TextChunk { + text_field = "text" + output_field = "chunk" + chunk_index_field = "chunk_seq" + chunk_size = 800 + overlap_size = 100 + } + Embedding { + model_provider = OPENAI + api_key = "sk-xxxx" + model = "text-embedding-3-small" + dimension = 1024 + single_vectorized_input_number = 10 + model_retry_max_attempts = 5 + vectorization_fields { + chunk = "chunk_vector" + } + } +} + +sink { + Milvus { + url = "http://127.0.0.1:19530" + token = "root:Milvus" + collection = "knowledge_base" + enable_upsert = false + batch_size = 500 + } +} +``` + +每个阶段发生的事情: + +1. **解析。** LocalFile source 把每个 Markdown 文件解析成元素行(标题、段落、列表项、代码块等),每行携带 `text`、`element_id`、`element_type`、`heading_level`、`position_index` 等字段。开启 `markdown_rag_metadata_enabled = true` 后,source 还会附加 RAG 元数据列:`source_uri`、`document_id`、`chunk_id`、`chunk_index`、`content_hash` —— 由文件路径与内容推导出的稳定标识。 +2. **切分。** TextChunk transform 把每个元素的 `text` 切成不超过 800 字符的文本块,相邻块之间保留约 100 字符的重叠。每个输出行保留全部原始字段,并追加块文本(`chunk`)及其在源行中的序号(`chunk_seq`)。 +3. **向量化。** Embedding transform 调用向量模型,把 `chunk` 转换为 1024 维向量,写入新字段 `chunk_vector`。 +4. **写入。** Milvus sink 在目标集合不存在时按上游 schema 自动建集合,并按每批 500 条写入。 + +> 注意:Markdown 解析器产出的 `chunk_index` 是从 1 开始的,而 TextChunk 自带的序号字段(`chunk_index_field`)是从 0 开始的。示例里把 TextChunk 的字段改名为 `chunk_seq`,避免两者冲突。 + +如果要摄入 PDF 文件,把 `file_format_type` 设为 `"pdf"` 并开启 `pdf_rag_metadata_enabled = true` 即可。 + +## 示例二:把数据库行向量化后写入 Qdrant + +同一套管道模式也适用于结构化数据 —— 这里把商品描述向量化后写入 Qdrant: + +```hocon +env { + parallelism = 4 + job.mode = "BATCH" +} + +source { + MySQL { + url = "jdbc:mysql://127.0.0.1:3306/shop" + driver = "com.mysql.cj.jdbc.Driver" + username = "root" + password = "******" + table_path = "shop.products" + } +} + +transform { + Embedding { + model_provider = OPENAI + api_key = "sk-xxxx" + model = "text-embedding-3-small" + dimension = 1024 + vectorization_fields { + description = "description_vector" + } + } +} + +sink { + Qdrant { + collection_name = "product_vectors" + host = "127.0.0.1" + port = 6334 + } +} +``` + +## 让向量保持新鲜 + +知识库是活的:文档会被编辑,数据库的行会变化。SeaTunnel 支持几种刷新策略: + +- **全量重建。** 把 Milvus sink 的 `data_save_mode` 设为 `DROP_DATA` 后重跑作业,丢弃旧集合并全量重建。最简单、永远一致,代价是重新跑一遍向量化。 +- **幂等更新。** 如果上游 schema 带有主键,把 Milvus sink 的 `enable_upsert` 设为 `true`,重跑作业时将按主键更新已有记录,而不是重复写入。 +- **变更数据捕获。** 使用 MySQL-CDC source 持续把行变更应用到向量库。CDC 需要搭配支持 upsert 的 sink,这样被更新的行会重新向量化并原地覆盖。 +- **块级生命周期。** 注意:当源文档发生变化时,SeaTunnel 不会自动删除过期的文本块(没有墓碑机制)。如果重新切分后块的数量变少了,旧块会残留。要么定期全量重建,要么自行设计稳定的块标识,并在管道的后续步骤里处理删除。 + +## 最佳实践 + +- **块大小。** `chunk_size` 按 Unicode 码点计数,默认 1000。`overlap_size` 要远小于 `chunk_size` —— 块数量大约按 `chunk_size / (chunk_size - overlap_size)` 增长。善用 `separators` 避免把句子从中间切断。 +- **向量化的可靠性。** 远程向量化 API 有限流,偶发失败。把 `model_retry_max_attempts` 设为大于 1,让限流、超时这类可重试错误按退避策略自动重试,而不是让作业直接失败。在服务商允许时,用 `single_vectorized_input_number` 把多条输入合并进一次请求。 +- **精度。** 向量目前只产出 `float32`,其他精度会被自动转换。务必保证 Embedding transform 配置的向量维度与目标集合期望的向量字段维度一致。 +- **主键。** TextChunk 要求 `output_field` / `chunk_index_field` 不能与主键或唯一键涉及的列重名。向 Milvus 写入 upsert 时,记住 upsert 要求集合 schema 中有主键。 +- **元数据过滤。** 开启 `markdown_rag_metadata_enabled = true` 后,每行都带有 `document_id`、`chunk_id` 和 `content_hash`。可以用 Metadata transform 把逻辑知识同步字段(如 `SourceUri`、`DocumentId`、`ChunkHash`)投影成应用自己的列,便于检索时过滤: + + ```hocon + Metadata { + metadata_fields = { + SourceUri = "ks_source_uri" + DocumentId = "ks_document_id" + ChunkHash = "chunk_hash" + } + } + ``` + +## FAQ + +### 文档变化时,SeaTunnel 会删除过期的文本块吗? + +不会自动删除。SeaTunnel 负责数据的搬运与转换,块删除策略属于消费方管道的职责。支持的策略见[让向量保持新鲜](#让向量保持新鲜)。 + +### 支持哪些向量模型服务商? + +支持 `OPENAI`、`AMAZON`、`DOUBAO`、`QIANFAN` 等常见服务商,另外提供 `CUSTOM` 模式 —— 自行定义请求头、请求体和响应解析。完整列表与选项见 [Embedding](../transforms/embedding.md) 文档。 + +### 一个文档变化后,必须把所有内容重新向量化一遍吗? + +不需要 —— 开启 `markdown_rag_metadata_enabled = true` 后,每个文本块都带有 `content_hash` 和稳定标识,下游系统可以据此判断哪些块发生了变化。注意:如果 source 与 sink 之间的 transform 改变了文本或将一行展开成多行(例如 TextChunk),最终的块标识和哈希应当在该 transform 之后重新计算。
