按真实协议配置集成。

与 FlowPort 0.1.5 对齐的云服务、分析数据库、服务发现与遥测边界。

目录与配置

提供 41 类原生 Go Source、34 类原生 Go Destination,另有 2 个 Source 与 8 个 Destination 服务预设。预设复用 Kafka、S3、Remote Write,不代表新增传输协议。安装与运行不依赖 Node.js 或 Vector 进程。

在 UI 选择服务卡片后填写真实地址、现有资源、认证与专用参数。复杂规则使用高级 JSON;保存草稿不会改变运行版本,修改后需重新发布。证书路径和环境身份属于实际执行 Worker。

服务预设

CLS、Azure Event Hubs 来源复用 Kafka 消费;SLS、CLS、Event Hubs 输出复用 Kafka。使用厂商端点、Topic、消费组及 SASL_SSL/PLAIN 凭证。SLS 输出使用 Project 用户名、AccessKey ID#Secret 密码和 Logstore.json Topic;Event Hubs 用户名为 $ConnectionString,密码为命名空间连接字符串。Kafka 预设按服务要求关闭幂等生产者,使用 gzip;Event Hubs 消费提交间隔默认 600ms。

R2、MinIO 输出复用 S3。填写实际 S3 Endpoint、Bucket、Access Key 和 Secret;R2 region=auto,MinIO 默认 us-east-1,预设启用 Path Style。不能只填写产品名称就认为协议与权限已验证。

VictoriaMetrics 单机使用 /api/v1/write,集群 vminsert 使用 /insert/<accountID[:projectID]>/prometheus/api/v1/write;Mimir 使用 /api/v1/push,tenant_id 对应 X-Scope-OrgID。使用完整写入 URL,租户不能代替网关认证。

SLS 原生 Source

aliyun_sls 设置 endpoint、project、logstore、group_id、access_key、secret_key;端点属于 SLS API,不是 Kafka 地址。预先创建 Project、Logstore 与消费组,offset 为 latest/earliest,默认每秒读取。

使用消费组心跳、分片检查点与接管;新分片共用一次检查点快照。消息经本地队列持久接受后才提交,队列拒绝或租约失效时不确认,重启/接管可能重复。

CloudWatch Logs

Source aws_cloudwatch_logs 设置 region、log_group_name,可选 log_stream_name、filter_pattern、start_time(RFC3339)。默认 latest 固定首次起点;earliest 用于补采。每页新事件与水位原子保存,重启从已接受水位恢复。

lookback_seconds 默认 60,范围 1–86400;窗口内按事件 ID 去重,最多保留 100000 个 ID。相同毫秒不同 ID 仍接收,窗口外迟到事件需另行补采。重复页不入队,扫描结束保存进度;空页不代表分页结束。

Destination 使用 PutLogEvents,设置 region、现有 log_group_name/log_stream_name,仅接收日志,不自动创建资源。Source 读取与 Destination 写入需要不同 IAM 权限,诊断通过不等于写权限通过。

Kinesis Data Streams

aws_kinesis Source/Destination 设置 region、stream_name。它与 Firehose HTTP 接入是不同协议。来源发现分片,先排空父分片再读取子分片,支持迭代器恢复、KPL 校验展开与持久序列号;latest 的首次时间边界会保存。一个流由一个 Worker 独占。

KPL 内部记录全部落盘后才推进物理序列号;中途失败可能重复已接受的记录。保留 Stream、Shard、Sequence、Subsequence、Partition Key 标签。

输出使用 PutRecords,单批最多 500 条/5 MiB,本实现每条连分区键最多 1 MiB。partition_key 或 partition_key_tag 提供分区键,本实现限制 1–256 字节。调用内仅重试临时失败记录;队列重试仍可能重复成功记录,不保证全局顺序或恰好一次。

RocketMQ

需要预建 Topic、消费组和 RocketMQ 5.x gRPC Proxy,endpoint 为 host:port;4.x Remoting 不属于该适配器。Source 使用 SimpleConsumer,配置 group_id、topic、filter 和 invisible_seconds(默认 120,范围 30–43200)。

处理期间续期不可见时间,落盘后 ACK;容量拒绝、租约失效或错误不确认。MESSAGE_NOT_FOUND 是空长轮询,保留连接并等待 200ms。Destination 同步发送并检查收据,每条消息最多 4 MiB。真实 RocketMQ 5.3.3 已验证空闲后发送、消费、持久接受与发送收据。

Doris 与 StarRocks

Destination 选择 doris/starrocks,设置 url、database、table,按服务填写 Basic 认证。使用 Stream Load JSON,支持 columns、jsonpaths、where、timezone,及 gzip。

队列批次 ID 与数据生成稳定 label,检查业务状态、加载行数和过滤行数;重复 label 仅在原作业 FINISHED 时确认成功。部分写入是永久失败,已写入行不自动重发。

307/308 只允许同 scheme、无嵌入凭证且主机在 redirect_hosts 中的跳转;配置完整 host:port。不会把认证头发送到未授权主机。诊断不导入业务记录,仍需真实表结构与写入联调。

GCS 与 Azure 对象通知

GCS Source 设置 project、subscription、bucket,经现有 Pub/Sub 订阅读取 OBJECT_FINALIZE;只下载配置 Bucket 的指定 generation,持久接受后确认。Worker 使用默认 GCP 身份或 credentials_json,需独立配置订阅与对象读取权限。

Azure Source 设置 url、container、queue,以及 servicebus_namespace 或 connection_string。Event Grid BlobCreated 通知进入预建 Service Bus 队列,只下载配置账户/容器且匹配 ETag 的对象;处理时续锁,落盘后完成消息。Blob 默认 Entra ID,也支持独立 Shared Key/SAS 字段;Blob SAS 不替代 Service Bus 身份。

支持 gzip、NDJSON/文本分批和有上限的 JSON;JSON 解压后最多 32 MiB,NDJSON 单行最多 1 MiB。不会扫描本地目录或全部 Bucket。GCS/Azure Destination 则是对象写入,通知读取与写入角色不同。

Prometheus 服务发现

静态 url/urls、HTTP SD、DNS SRV/A/AAAA,以及远端 Kubernetes EndpointSlice/Consul passing 服务可组合。配置 kubernetes_sd_url 或 consul_sd_url,并使用独立发现凭证。HTTP SD 与抓取共用认证/TLS,应使用可信发现服务。

高级 JSON 支持 target_relabel_configs/metric_relabel_configs,各最多 64 条,以及每行一个、匹配完整指标名的 metric_include/metric_exclude。每个运行版本只解析一次规则;修改后重新发布。

按 relabel 后的 URL 和普通标签共同去重;同 URL 不同标签独立采集,相同目标只抓取一次。max_targets 默认 256(1–1000),concurrency 默认 4(1–32)。失败保留进程内前次目标,成功空列表移除;缓存不跨重启。

协议保真与输出边界

OTLP HTTP/gRPC 保留 Resource、Scope、schema、Sum temporality/monotonic、Histogram、ExponentialHistogram、Summary、exemplar,以及 Span Kind、状态、事件和 Links。Pipeline 数值/标签修改会合并到输出,过滤 Point 同时移除原始协议记录。

Remote Write v1 来源保留 native histogram/exemplar;普通 Point 或跨协议来源不自动转换成等价的原生直方图,不发送 metadata 或 Remote Write v2。普通数值样本仍使用 float64 和毫秒,不能保留任意大整数或纳秒精度。

Vector gRPC 使用 v2 PushEvents;普通 Point 指标映射 Gauge,不恢复 Counter/Histogram 或原始 Vector envelope。内部 __flowport_ 字段在 JSON、文本与 Line Protocol 输出隐藏。OTLP 部分拒收不作为全部成功或自动重发;以实际服务收据判断交付。

诊断与验证范围

在实际 Worker 测试连接和试采样。连接诊断不替代写权限、表结构或业务交付验证;消息/推送和持续云来源从运行输入旁路采样,不创建额外消费者或推进游标。

当前版本通过 Go 测试、race、协议模拟、真实 RocketMQ 收据测试、五语言检查和两架构无 Node.js 的离线安装运行验收。真实云 IAM、配额、网络和长稳仍需在客户环境联调。社区版一个本机 Worker;多 Worker 需要离线 License。