Solutions · Open Source

Argus Flow

将 Apache NiFi 2.10.0 部署到生产环境所需的一切集中于一处的开源数据管道发行版。自定义扩展包(NAR)、发布包(tar.gz、RPM、容器镜像)、Kubernetes Operator,以及配置、证书与用户管理工具——把原本分散的仓库整合为单一 monorepo,由同一条构建流水线产出版本相互对齐的产物。

Apache License 2.0 · 开源Apache NiFi 2.10.0 · Java 21GitHub 仓库

特点与优势

01

扩展、发布与 Operator 同一条流水线

NiFi 扩展、Operator 与证书生成脚本此前分散在各自仓库,版本容易错位。现已整合为 monorepo:一个 make 目标即可产出版本对齐的 NAR、tar.gz、RPM 与容器镜像。

02

12 个包 · 18 个 NAR 实现依赖隔离

扩展按领域拆分为独立包,每个包单独打成 NAR,第三方依赖不会跨越包边界发生冲突。可在 CDC、湖仓、Hive、Parquet、数据库、报告与流程分析中按需安装。

03

无需 ZooKeeper 的 NiFi 2.x Kubernetes Operator

直接使用 NiFi 2.x 自带的 Raft 协调机制,无需 ZooKeeper 即可运行集群。通过 NiFiCluster、NiFiFlow 自定义资源进行声明式管理,缩容时按 Offload → Disconnect → Remove 顺序安全下线,并通过 cert-manager 自动签发 TLS 证书。

04

用工具安全管理配置、证书与账号

不必手工编辑超过 200 行的 nifi.properties 与 XML,交互式向导收集取值并在保存前展示 diff。--check 诊断可发现证书 SAN 不一致、临近过期以及被注释掉的登录 Provider,避免 Invalid SNI 之类的故障。

发行版构成

扩展包、发布包、Kubernetes Operator 与运维工具在同一个仓库中一起构建,同时每个目录也可以独立构建。

NiFi Extensions
Maven · Java 21
12 个包 → 18 个 NAR(按包隔离依赖)
处理器 22 种 · 控制器服务 7 种
报告任务 5 种 · Flow Analysis Rule 6 种
基于数据库的认证与鉴权 Provider
以 services-api 包拆分服务契约
统一包名 io.datadynamics.nifi.*
Distribution
tar.gz · RPM · 容器
对官方 NiFi 二进制包重新打包
基于 nfpm 的 RPM 与 systemd 单元
内置 NAR 的容器镜像(apache/nifi 基础镜像)
bin/ 目录随附运维脚本
上游不做 vendoring,构建时锁定版本
make dist · make rpm · make docker-image
Kubernetes Operator
Python · kopf
NiFiCluster、NiFiFlow 自定义资源
协调 StatefulSet、Headless Service、ConfigMap
扩缩容时安全下线节点
通过 cert-manager 自动签发与续期 TLS
基于健康检查自动重连节点
Prometheus 指标 · Helm Chart
Ops Tooling
argus-config · argus-ssl · argus-user
conf/ 交互式配置 TUI(以 zipapp 随发行版分发)
初次配置向导 —— 访问地址 → TLS → 登录方式
--check 配置诊断(退出码可接入 CI)
非交互式 --set、--recipe 自动化
TLS 证书生成脚本
数据库认证的用户管理 CLI
技术栈
Apache NiFi 2.10.0Java 21Maven 3.9.16+ (wrapper)Python 3.12+kopfHelmnfpmDocker / Podmancert-managerDebezium 3.2.2Apache IcebergDelta KernelApache KuduHive3 · ORCParquet · AvroPostgreSQL

核心能力

从 CDC、湖仓、Hive 与文件格式扩展,到运维治理规则、数据库认证鉴权、发布打包、Kubernetes Operator 以及配置与 TLS 自动化——把 NiFi 运维所需的 12 个支柱汇入一个发行版。

4 种 CDC 数据源

基于 Debezium Embedded Engine 3.2.2,将主流关系型数据库的变更事件输出为 JSON FlowFile。

CaptureChangeMySQL —— binlog CDC
CaptureChangePostgreSQL —— 逻辑复制(pgoutput/decoderbufs)
CaptureChangeOracle —— LogMiner(CDB/PDB)
CaptureChangeSQLServer —— SQL Server CDC
PrimaryNodeOnly · at-least-once

湖仓写入

将记录直接写入开放表格式与列式存储。

PutIceberg —— 写入 Apache Iceberg 表
PutDeltaLake —— 基于 Delta Kernel 的路径型写入
PutKudu —— 写入 Apache Kudu 表
HiveCatalogService · HadoopCatalogService

数据库处理器

从查询到批量写入,拓宽了关系型数据库的对接路径。

ExecuteSQL · ExecuteSQLRecord
ExecuteFastSQL —— 大结果集高速流式读取
PutDatabaseRecord —— INSERT/UPSERT
BulkOracleInsertProcessor —— Oracle 数组绑定批量插入

Hive / ORC

以专用包提供 Hive3 对接与 ORC 写出。

SelectHive3QL · PutHive3QL · PutHive3Streaming
UpdateHive3Table —— 列与分区的表结构同步
TriggerHiveMetaStoreEvent
PutORC · Hive3ConnectionPool · VendorHive3ConnectionPool

文件格式与记录序列化

提供可控制时间戳格式的读取器与写入器。

PutParquet · MergeParquet
CSVReader —— 可指定时间戳格式
TimestampFormatAvroRecordSetWriter
TimestampFormatParquetRecordSetWriter

通用处理器

把现场反复需要的处理收进标准包。

PrePostExecuteStreamCommand —— 附带前置/后置命令的执行
MultilineCsvParser —— 解析含多行字段的 CSV
按包分发,仅安装所需部分

5 种监控报告任务

按阈值监控节点资源状态并向外部告警。

MonitorDiskUsageReportingTask —— 仓库磁盘
MonitorMemoryUsageReportingTask —— JVM 堆
MonitorMemoryPoolReportingTask —— 单个内存池
MonitorThreadReportingTask —— 死锁等线程状态
HttpNotificationReportingTask —— Webhook HTTP POST 通知

6 种 Flow Analysis Rule

以静态检查发现流程配置反模式的运维治理规则。

TimerThreadPoolCeilingRule —— 全局线程数上限
ProcessorThreadShareLimitRule —— 单个处理器占用比例
ProcessorConcurrencyCapRule —— 按类型的并发任务上限
ListingScheduleGuardRule —— List 类处理器的 0 秒调度
DeadEndFunnelRule —— 死端 Funnel 造成的积压
IcebergSinkMergeRule —— 缺少 Merge 导致的小文件问题

基于数据库的认证与鉴权

为既无法引入 LDAP 也无法引入外部 IdP 的环境准备的 Provider,避免手工修改 users.xml。

DbLoginIdentityProvider —— 从关系库读取用户与口令进行认证
DbUserGroupProvider —— 基于关系库的鉴权
argus-user.sh 用户管理 CLI
默认关闭(注释状态),按需启用
内置 PostgreSQL 驱动,MariaDB 需另行提供

发布打包

重新打包官方二进制,产出可直接投入生产的交付物。

tar.gz —— 包含扩展 NAR 与运维脚本
RPM(nfpm)—— 随附 systemd 单元与 sysconfig
容器镜像 —— 基于官方 apache/nifi
上游不做 vendoring,构建时锁定版本

Kubernetes Operator

以声明式方式定义 NiFi 2.x 集群,并自动完成扩缩容与更新。

NiFiCluster CR —— 节点数、资源、存储与 JVM 堆
NiFiFlow CR —— 以 GitOps 方式部署 Registry Flow
通过 Headless Service DNS 自动发现节点
滚动重启实现配置的无中断生效
对接 cert-manager Issuer 的 TLS
Helm Chart —— CRD、RBAC、Deployment

配置与 TLS 自动化

把安装之后最容易出问题的配置环节用工具标准化。

初次配置将访问地址原样写入证书 SAN
--check 诊断 SAN 不一致、证书过期与被注释的 Provider
配方 tls:generate · login:db · login:ldap · state:zookeeper
保留原有注释与顺序,保存前展示 diff
非交互模式只打印证书命令而不执行

版本

以 Community 自由使用整个开源发行版;当需要运维 SLA 与专属技术支持时,可选择 Enterprise。

Community

Apache License 2.0 · 免费

不受限制地使用并自行运维扩展包、发布包、Operator 与运维工具的全部内容。

推荐

Enterprise

企业客户支持

在 Community 全部内容之上,增加基于 SLA 的专属技术支持,以及实施、迁移与定制扩展开发。

功能对比
Community
Enterprise
核心功能
NiFi 扩展 12 个包 · 18 个 NAR
CDC湖仓Hive / ORCParquetDatabase报告Flow Analysis
发布包(tar.gz · RPM · 容器镜像)
Kubernetes Operator · Helm Chart
配置 TUI · 诊断 · TLS 证书工具
基于数据库的认证与鉴权 Provider
技术支持与服务(Enterprise)
基于 SLA 的专属技术支持
优先提供热修复与安全补丁
安装、实施与 NiFi 1.x → 2.x 迁移支持
培训、上手指导与管道架构咨询
定制处理器开发与路线图优先响应
隔离网络部署与厂商 JDBC 驱动对接支持
Cloudera Hive JDBCMariaDB Connector/J
支持渠道
支持渠道
GitHub Issues
专属支持渠道
Apache License 2.0 · 开源

以开源方式发布的 NiFi 生产发行版

Argus Flow 以 Apache License 2.0 在 GitHub 上公开。从扩展包源码到发布打包、Kubernetes Operator 与运维工具全部开源,企业可以自行审计代码、按自身环境扩展,并在隔离网络中完全自主运行。

  • 商用无限制的 Apache 2.0
  • 上游归属与许可证头在构建中自动校验
  • 本地化与隔离网络自主运行