9. ExecuteFlinkSQL 处理器参数设计
ExecuteFlinkSQL处理器参数包含以下三类:
1. flowfile 读写参数
为了从flowfile 中读取数据,并转换为数据行,我们可以利用 NiFi 的 Record Reader服务来完成;为了将数据输出,需要将数据转换为flowfile,我们可以利用 NiFi 的 Record Writer 服务来完成。因此我们设计了以下两个参数:
2. FlinkSQL 程序参数
为了从 NiFi 接收数据流,需要在 Flink 中创建一个表数据源,我们可以在 Flink 中执行一条 SQL 来实现。例如通过以下 Flink SQL 语句:
CREATE TEMPORARY TABLE tab_nifi ( `name` STRING NOT NULL, `proc_time` AS PROCTIME()) WITH ( 'connector'='nifi')
这里必须指定 ’connector’=’nifi’,来创建了一个数据源表,从 NiFi 接收数据。
另外,我们还需要设定一条查询语句,实现对数据的处理,例如:
SELECt window_start, window_end, name, COUNT(*) AS `count` FROM TABLE(TUMBLE(TABLE tab_nifi, DESCRIPTOR(proc_time), INTERVAL '1' MINUTES)) GROUP BY window_start, window_end, name
这样我们就利用一个滚动窗口,按分钟统计了数据条数。这两个参数配置图如下:
3. NiFi 与 Flink 通信参数
我们使用 HTTP 实现 NiFi 与 Flink 间通信,我们需要指定侦听接口、侦听端口,希望使用 https 时可以设置 SSL 上下文服务,客户端认证方式,以及请求队列长度和请求超时时间等参数:
4. Flink 执行参数
Flink 执行参数包括并行度,以及使用外部 Flink 集群时,需要指定 Flink JobManager 的主机地址和端口:
10. 代码工程设计
对于代码工程,我们可以从 NiFi 插件的样例代码开始。一般情况下,一个NiFi插件需要两个工程:
1. 插件打包工程(nifi-flink-sql-nar);
2. 处理器/服务代码工程(nifi-flink-sql-processors);
对我们来说,要把nifi-connector注入到flink环境中执行,最好切分出一个单独的工程(nifi-flink-sql-connectors)来。所以我们总共有三个工程,代码在 processors 和 connectors 两个工程中的情况如下:
11. 测试和运行
我们使用 GenerateFlowFile 对 ExecuteFlinkSQL 处理器进行测试:
从 Flink Web 界面上,也可以看到任务正在运行:



