企业服务与数据采集
FlowPort 0.1.6 增加 10 类原生 Go Source、5 类 Destination,以及 Grafana Cloud / New Relic 两个 OTLP 服务预设。目录合计 51 类 Source、39 类 Destination,另有 2 个来源与 10 个输出预设。安装与运行仍不需要 Node.js 或 Vector。
Azure Service Bus
来源与目标类型均为 azure_service_bus。填写 servicebus_namespace,例如 example.servicebus.windows.net;认证使用 Worker 的 Entra ID,或独立 connection_string。来源二选一:queue,或 topic + subscription;目标选择一个 queue 或 topic。资源需要预先创建。
来源使用 PeekLock,每次一条消息,处理时每 5 秒续锁,本地持久接受后才 Complete。容量拒绝、租约失效或取消不确认。发送使用 SDK 协商容量的消息批次并等待发送收据,Message ID 根据持久批次身份和行内容稳定生成;仅在上游实体开启重复检测时有相应去重效果,不保证端到端恰好一次。部分批次成功后整分支重试可能重复。
连接诊断不消费或发送业务消息:来源 Peek,目标建立发送链路和空消息批次。读取或链路权限不等于完整交付验证。试采样旁路观察运行输入。
SQL 查询与增量
来源为 postgresql_query、mysql_query、sqlserver_query。url 不嵌入凭证,用户名、密码放在独立字段。示例地址为 postgres://HOST:5432/DB、mysql://HOST:3306/DB、sqlserver://HOST:1433?database=DB。默认 30 秒查询一次;一页 1000 行,可设 1–10000。
仅允许单条 SELECT,不接受 SELECT INTO 或多语句。PostgreSQL/MySQL 使用只读事务;SQL Server 必须使用仅有 SELECT 权限的账号。没有 cursor_column 时,每轮完整查询,结果可重复。查询参数使用数据库驱动绑定,不能拼接游标值到 SQL。
增量配置 cursor_column、cursor_start,查询包含 {{cursor}} 与 {{limit}},并按唯一游标升序排列。同一时间有多行时配置 tie_breaker_column / tie_breaker_start,查询使用 {{tie_breaker}},按时间与唯一主键共同排序。每页数据与下一游标原子入队;拒绝入队不推进,重启从持久位置恢复。新查询、资源或游标列形成新进度空间。
SQL Server 使用 SELECT TOP ({{limit}}) ... WHERE id > {{cursor}} ORDER BY id。定时增量查询无法捕获物理删除,也无法保证捕获未更新水位的改动;这些需求使用 CDC。
SELECT id, message FROM events
WHERE id > {{cursor}} ORDER BY id LIMIT {{limit}}PostgreSQL CDC
postgresql_cdc 需要 url、预建 slot、publication、数据库身份及复制权限;服务端开启 wal_level=logical,复制槽使用 pgoutput。需要完整更新/删除旧值时,为表配置 REPLICA IDENTITY FULL;默认身份只提供可用的键值。未改变的大 TOAST 值使用 unchanged_toast 标记,不能当作 null。
默认从复制槽确认位置或 start_lsn 读取;已有持久进度优先。按事务收集 insert/update/delete/truncate,提交后将事务与结束 LSN 原子入队。表结构消息更新列映射,不自动变更目标表。仅所有目标完成后的可交接位置向复制槽确认,避免新 Worker 跳过仍在投递的批次。复制槽会保留 WAL,需监控磁盘与复制滞后。
开启 initial_snapshot 时创建临时逻辑槽,导出一致性快照,读取 publication 的列、行过滤与分区范围,再从快照对应 LSN 进入变更流。用户的持久复制槽不自动创建或删除。PG 快照需要 PostgreSQL 15+ 及支持导出快照的复制权限。
MySQL CDC
mysql_cdc 使用原生 Go binlog 客户端,需要 MySQL 8.0+、唯一 server_id、复制权限,及 binlog_format=ROW、binlog_row_image=FULL、binlog_row_metadata=FULL。配置 tables 为每行 database.table;留空时读取全部行事件。支持文件/位置及可选 gtid_set,已有进度优先,未指定起点时固定并保存当前 binlog 位置。
只在事务提交后持久接受数据及进度;失败不会保存新位置。支持增删改与选定表相关的 DDL 事件,刷新后的列元数据用于后续行。DDL 输出不执行在下游;不接受语句模式数据或 JSON 局部更新编码。保留足够 binlog,不会自动修复已被上游清除的历史。
初始快照须指定 tables,且全部为 InnoDB。建立快照时使用短暂全局读锁,固定一致性读视图与 binlog 位置后立即解锁;需要 FLUSH_TABLES/RELOAD 等服务端权限,可能短暂阻塞写入。随后分批读取快照并衔接 binlog。
两类 CDC 的事务缓冲默认 32 MiB,可设 1 KiB–128 MiB;超限停止,不确认或静默拆分事务。快照每页默认 500 行。输出包含 operation、schema/table、before/after 和协议进度;快照另有 snapshot_begin、snapshot、snapshot_complete 与 snapshot_id。快照中断后创建新的快照标识并重新读取,可能重复已接受行;下游物化时须按快照标识和完成边界协调,不能宣称端到端恰好一次。试采样旁路观察运行输入,不创建复制槽或读取额外变更流。
{"type":"postgresql_cdc","config":{"url":"postgres://HOST:5432/DB","username":"replication_reader","password":"REPLACE_ME","slot":"flowport_slot","publication":"flowport_pub","initial_snapshot":true,"max_transaction_bytes":33554432}}Microsoft 365 与 Google Workspace 审计
microsoft365_audit 配置 tenant_id、client_id、client_secret,使用 OAuth 客户端凭证;可使用短期 token,到期需更新。content_type 默认 Audit.General,可按已授权内容类型调整。先开启统一审计与内容订阅;只有 create_subscription=true 时,运行来源才调用订阅创建,诊断不会创建订阅。
google_workspace_audit 配置开启域委派的 credentials_json、有审计读取权限的 delegated_user、application_name(默认 login)、user_key(默认 all)。授予 admin.reports.audit.readonly 范围;可选短期 Token 不自动续期。委派服务账号的 Token 自动更新。
两类来源默认每 60 秒拉取,初始读取回看一小时。可填 RFC3339 start_time,每个窗口最多 24 小时,并保留 2 分钟延迟。lookback_seconds 可设 1–86400;按内容/事件 ID 去重,窗口最多 100000 个 ID,单轮最多 100 页。分页和数据一起落盘,重启继续;下一页与内容地址必须同域,避免泄露认证。实际历史范围、事件迟到、权限、许可及配额以云服务为准;回看窗口外迟到数据需补采。试采样使用运行输入旁路。
SNMP 与 SNMP Trap
snmp 填写远端 address=host:161 与每行一个数字 oids,默认 SNMPv2c / community;支持 SNMPv1,以及 SHA256 + AES128 的 SNMPv3 authPriv。v3 填写 username、auth_password、privacy_password。按 OID GET,walk=true 遍历子树;每轮变量默认最多 1000,可设 1–10000。保留 OID、主机及值;非数值变量不自动转换为数值指标。
snmp_trap 监听 UDP,默认 127.0.0.1:1162;上游发送到可达地址。认证版本与 community/用户必须匹配。v3 还应配置适当的十六进制 engine_id。普通 Trap 没有确认或可靠重放;Inform 仅在持久接受后响应,拒绝入队不成功确认。Community 与密码在 API/UI 遮蔽,不写入事件字段。
Azure Monitor Logs
目标 azure_monitor_logs 填写 DCR/DCE 的 HTTPS 基地址 url、DCR 的 immutable dcr_id 和输入 stream_name。表与 DCR 预先创建,授权身份拥有相应发布权限。使用 Worker Entra ID,或 tenant_id / client_id / client_secret,也可填短期认证 Token。
按 JSON 数组分批,每个请求最多 512 KiB。默认发送 measurement/tags/fields/timestamp/category;DCR 的输入声明须匹配,可在高级 JSON 配置 field_mapping 将目标列映射到 fields.*、tags.*、measurement、timestamp、category,再由 DCR 转换。认证探测不会写入数据,不能证明 DCR 写权限。
BigQuery Storage Write
目标 gcp_bigquery 填写 project、dataset、table,使用 Worker ADC、服务账号 JSON 或短期 Token。通过 Storage Write API _default 流写入,读取真实表结构并编码 protobuf,检查每次 AppendRows 的响应;每请求最多 4 MiB。
默认表包含 measurement STRING、tags JSON、fields JSON、timestamp INT64(纳秒)、category STRING。高级 JSON 的 field_mapping 可映射到实际列;支持 STRING、JSON、INT64、DOUBLE、BOOL、BYTES 与 TIMESTAMP,TIMESTAMP 转为微秒;不支持复杂/重复列时明确失败。JSON 字段保留大整数文本。默认流至少一次,不提供跨重试 offset 去重;需要下游按稳定业务标识去重。行错误不记为成功,不自动建表。诊断仅读取结构。
AWS Data Firehose
目标 aws_firehose 填写 region、stream_name,使用 Worker 凭证链或 Access Key/Secret/Session Token。与 Firehose HTTP 输入及 Kinesis Data Streams 是不同方向和协议。
PutRecordBatch 每批最多 500 条/4 MiB,单条最多 1000 KiB。检查逐条收据,仅重试暂时失败记录,最多五轮;整队列分支重试仍可能重复已成功记录。成功表示 Firehose 接受,不表示其最终目标已经写入。诊断使用 DescribeDeliveryStream,不发送记录。
Snowflake Streaming
目标 snowflake 配置 HTTPS url、database、schema、table。token_type=scoped 时 url 为 ingest 基地址,认证字段填写现有 Scoped Token;oauth 时 url 为账户地址,填写现有 OAuth Token;jwt 时配置 account、username、至少 2048 位的 RSA private_key PEM,由 Go 签发 JWT、发现 ingest 地址并兑换 Scoped Token。私有接入可显式填写可信 HTTPS ingest_url,仍验证证书。Token/私钥在配置中加密并通过 API 遮蔽。
使用 Elastic Channel REST,分批 NDJSON,每批原始数据最多 3 MiB,可选 gzip;高级 field_mapping 映射列。请求 ID 从持久批次和内容稳定生成,HTTP 成功后仍检查 message=OK。表示服务持久缓冲,不表示目标表立即可查或每行转换成功;在 Snowflake 检查错误表。Elastic 为至少一次且无顺序保证,模糊收据后重试可能重复,不宣称 Named Channel 的恰好一次。Scoped 模式没有安全的写入探测,诊断明确保留未知范围。
Grafana Cloud 与 New Relic 预设
两者复用现有 OTLP HTTP,支持 logging/metric/tracing,默认 gzip。Grafana Cloud 填控制台提供的 OTLP 基地址(保留 /otlp 前缀)与 Basic 用户名、Token 密码;New Relic 的 API Key 填 Ingest License Key,以 api-key 请求头发送,按账户区域修改地址。参数预设不代表对应账户已联调。
验证范围
OrbStack 中的真实 PostgreSQL 与 MySQL 验证 SQL 增量、初始快照、增删改和重启恢复。协议测试覆盖审计拒绝入队/去重、SNMP Inform 确认、Firehose 部分失败、BigQuery protobuf/拒收、Azure DCR 请求和 Snowflake 收据。SQL Server 与真实云账户的身份、配额、表结构、网络、服务版本和长稳仍需实际环境联调。连接诊断与 Pipeline 样本运行均不能替代真实端到端交付。