> ## Documentation Index
> Fetch the complete documentation index at: https://docs.molesignal.com/llms.txt
> Use this file to discover all available pages before exploring further.

# 调度流水线

> 按调度读取源数据流，用 VRL 转换链处理，再写入目标数据流或外部 connector。

**调度流水线**按时间窗读取源数据流，依次应用一组 VRL 转换，并把结果写入目标数据流（标准
intake）——还可把同一批事件扇出到外部 **connector**（S3、Kafka）。流水线按简化间隔计划运行，也可对
历史窗口按需做一次**回填（backfill）**。

<Note>
  调度流水线在数据落库**之后**作为周期任务处理。若想在 intake 热路径上即时转换事件，请改为在 intake
  步骤挂一个[函数](/zh-Hans/functions)。
</Note>

<Frame caption="调度流水线">
  <img src="https://mintcdn.com/molesignal/W03b-Z-TATDejvIA/images/architecture/pipeline_zh-Hans_light.svg?fit=max&auto=format&n=W03b-Z-TATDejvIA&q=85&s=3c6ed7fd6ae1062fb4576575178c3d8c" alt="调度流水线" className="block dark:hidden" width="477" height="595" data-path="images/architecture/pipeline_zh-Hans_light.svg" />

  <img src="https://mintcdn.com/molesignal/W03b-Z-TATDejvIA/images/architecture/pipeline_zh-Hans_dark.svg?fit=max&auto=format&n=W03b-Z-TATDejvIA&q=85&s=840feb2307731e9f1bbf150fdc622a22" alt="调度流水线" className="hidden dark:block" width="477" height="595" data-path="images/architecture/pipeline_zh-Hans_dark.svg" />
</Frame>

## 在图形编排器里搭建

打开 **流水线 → 新建**，在画布上排布 **来源 → 转换 → 目标**。节点可拖动，可在端点之间连线，编排器会随
时校验图（缺来源/目标、转换缺名称或脚本、流名重复、来源与目标同名等）。第一个来源和第一个数据流目标会
保存为流水线的 `source_stream` 与 `target_stream`。

| 节点     | 含义                             |
| ------ | ------------------------------ |
| **来源** | 每次运行读取的数据流。                    |
| **转换** | 一个 VRL 步骤。多个步骤自上而下按序执行。        |
| **目标** | 目标数据流（标准 intake）或外部 connector。 |

## 转换

每个转换是一段 [VRL](/zh-Hans/functions) 脚本。转换检视面板的「复用」下拉里，可选用
**内置预设**或任意已保存的[函数](/zh-Hans/functions)——选中即把脚本灌入该步骤，之后
可内联编辑。

内置预设随每个实例提供，且只读：

| 预设                 | 作用                                                |
| ------------------ | ------------------------------------------------- |
| `normalize-logs`   | 解析 JSON `message`，补充环境/集群字段，统一 `level` 小写。        |
| `route-by-service` | 归一化 `service` 并派生每服务目标流名（`logs_<service>`）。       |
| `parse-key-value`  | 把 `logfmt` / `key=value` 形式的 message 解析为字段后合并回事件。 |
| `redact-email`     | 把 `message` 中的邮箱地址替换为占位符。                         |
| `add-intake-time`  | 附加摄取时间戳，便于排查端到端延迟。                                |

<Tip>
  预设与已保存函数共用一个目录，二者都出现在**函数**页。内置预设带标识且只读；使用**复制脚本**可
  创建新的自定义函数。
</Tip>

### 扩展表

**扩展表**是一张 key → 记录的查找表，可在转换里 join。打开**扩展表**，新建表并添加行（一个 key 加若干
命名字段），再在 VRL 里查找：

```ruby theme={null}
# 用所属团队丰富每条事件
.team = lookup("service_meta", .service).team
```

查找在内存中完成，不会给运行增加查询开销。

## 目标与 connector egress

目标可以是**目标数据流**（默认——事件经标准 intake 写入并可被查询），也可以是外部 **connector**：

| Connector | 投递负载                                                      |
| --------- | --------------------------------------------------------- |
| **S3**    | 批量 JSON Lines `PUT` 到 `s3://<bucket>/<prefix><id>.jsonl`。 |
| **Kafka** | 每条事件产生一条 JSON 记录到配置的 topic。                               |

在**流水线 → Connectors**添加 connector，再在图形编排器里选择该 connector 作为目标。一条流水线可在同一次运行里既写
数据流又向 connector egress。

### 把源数据流从查询中隐藏

将一个兜底源数据流扇出成每服务（或每租户）目标流时——例如用上面的 `route-by-service` preset——源数据流
会继续堆积未分流的原始事件。要让查询和仪表盘只对准分流后的目标流，把源数据流标记为**不可查询**：打开
**数据流 → *源数据流* → 设置**，关闭 **可查询**。

不可查询的数据流照常采集并保留数据（也照常为本流水线供数），但会从查询选择器中隐藏，且 SQL 与 PromQL
搜索都会以 `stream is not queryable` 拒绝。随时可以再打开开关。

## 调度与回看窗口

| 字段              | 说明                                      |
| --------------- | --------------------------------------- |
| `cron`          | 间隔简写：`every:30s`、`every:5m`、`every:1h`。 |
| `lookback_secs` | 每次运行从源回看多久（默认 300）。                     |
| `enabled`       | 不删除流水线即可开关调度。                           |

Runner 不解析标准 cron 表达式；不符合 `every:` 语法的计划会被跳过。

每次触发，runner 读取源的 `[now - lookback_secs, now]`，应用转换链并写出结果。每次运行都会记录——在流
水线的 **运行** 标签页查看，或通过 `GET /api/v1/scheduled_pipelines/{id}/runs`（`state`、`scanned_rows`、
`error`）。

<Note>
  调度 runner 是单例：分布式集群下只在 **alert-manager**（或 **standalone**）节点运行，因此启用的流水线
  每个间隔只触发一次，而非每节点一次。
</Note>

## 回填（Backfill）

要按需处理历史窗口，提交一次回填：

```bash theme={null}
curl -X POST http://localhost:5080/api/v1/scheduled_pipelines/$PIPELINE_ID/backfill \
  -H "authorization: Bearer $MS_JWT" \
  -H 'content-type: application/json' \
  -d '{"start_micros": 1717200000000000, "end_micros": 1717286400000000}'
```

窗口必须 **≤ 31 天**。请求返回 `202` 带 `job_id` 和监控 URL；回填走异步
[search-job](/zh-Hans/distributed-deployment#async-search-job-pipeline) worker——读取源窗口、应用同一条
转换链、写目标流并 egress。

## 权限

* `pipelines.read`：列出流水线并读取运行历史。
* `pipelines.create`：创建流水线。
* `pipelines.edit`：修改图、转换、调度或回看窗口。
* `pipelines.pause`：启用或暂停调度执行。
* `pipelines.run`：提交回填任务。
* `pipelines.delete`：删除流水线。

<Card title="流水线 API" icon="code" href="/zh-Hans/api/intake/scheduled-pipelines">
  通过 HTTP API 创建、更新、列出与回填调度流水线。
</Card>
