惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

WordPress大学
WordPress大学
B
Blog RSS Feed
The GitHub Blog
The GitHub Blog
爱范儿
爱范儿
博客园 - 司徒正美
J
Java Code Geeks
酷 壳 – CoolShell
酷 壳 – CoolShell
Engineering at Meta
Engineering at Meta
大猫的无限游戏
大猫的无限游戏
D
Docker
Blog — PlanetScale
Blog — PlanetScale
Recent Announcements
Recent Announcements
罗磊的独立博客
Microsoft Azure Blog
Microsoft Azure Blog
博客园 - 聂微东
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
人人都是产品经理
人人都是产品经理
Stack Overflow Blog
Stack Overflow Blog
M
MIT News - Artificial intelligence
腾讯CDC
T
The Blog of Author Tim Ferriss
小众软件
小众软件
U
Unit 42
T
Tailwind CSS Blog

博客园 - 自由港

postgres 支持全文索引 查看 milvus 中的数据 查询数据库锁死情况 使用 keepalived 实现 tendis 高可用 部署tendis IDEA 运行 main 方法导致整个项目 install 自定义classloader hive 基础操作 在 服务器部署 seatunnel web服务 使用 seatunnel web 设计一个数据同步 如何 运行 seatunnel web 开发版 使用 seatunnel 实现数据同步 avro 数据入门 NIFI国际化 使用 NIFI读取EXCEL 数据到数据库 使用NIFI 同步数据库表 使用 NIFI监控数据库表 切换JDK NIFI实现配置分发
NIFI 使用HTTP 作为数据源接收数据
自由港 · 2025-11-07 · via 博客园 - 自由港

1.概述

在NIFI 中,可以 ListenHTTP 组件 启动一个HTTP服务,通过HTTP 服务接收 客户端 发送的信息,后续可以增加处理器,对请求进行处理。
我做了一个示例

  1. 通过 ListenHTTP 接收信息
  2. 通过 ExecuteGroovyScript 对数据进行处理
  3. 通过 PutFile 对数据进行验证,查看数据是否有效

2.配置

2.1 配置ListenHTTP 组件

image

配置后,我们 就可以通过 将数据 post 到 http://ip:8080/data 上进行数据接收

2.2 ExecuteGroovyScript 处理脚本

import groovy.json.JsonSlurper
import com.fasterxml.jackson.databind.ObjectMapper
import com.fasterxml.jackson.databind.json.JsonMapper
import com.fasterxml.jackson.core.JsonGenerator
import org.apache.commons.io.IOUtils
import java.nio.charset.StandardCharsets

def flowFile = session.get()
if (!flowFile) return

try {
    def text = ''
    session.read(flowFile, { inputStream ->
        text = IOUtils.toString(inputStream, StandardCharsets.UTF_8)
    } as InputStreamCallback)

    // 用 Groovy 解析 JSON(支持宽松语法)
    def jsonSlurper = new JsonSlurper()
    def data = jsonSlurper.parseText(text)
    data['age'] = 20

    // 使用 Jackson 序列化,不转义非 ASCII 字符(保留中文)
    def mapper = new ObjectMapper()
    mapper.getFactory().setCodec(mapper)
    // 关键:禁用 Unicode 转义
    def writer = mapper.writer()
        .without(JsonGenerator.Feature.ESCAPE_NON_ASCII)

    def updatedJson = writer.writeValueAsString(data)

    flowFile = session.write(flowFile, { outputStream ->
        outputStream.write(updatedJson.getBytes(StandardCharsets.UTF_8))
    } as OutputStreamCallback)

    flowFile = session.putAttribute(flowFile, "mime.type", "application/json; charset=utf-8")
    session.transfer(flowFile, REL_SUCCESS)

} catch (Exception e) {
    log.error("处理FlowFile时发生错误", e)
    session.transfer(flowFile, REL_FAILURE)
}

这个代码的作用是读取 上游的数据,并通过脚本修改这个数据,再投放给下游。

2.3 配置 PutFile

这里只需要配置一下执行路径就可以了

image

2.4 测试

我们使用 apipost 工具给 http 服务发一个请求

image

我们可以查看生成的文件内容

{"name":"张飞","age":20},我们可以操作这个JSON.