You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何修改NiFi Groovy脚本实现ClickHouse批量插入查询

修改后的NiFi Groovy批量生成ClickHouse插入语句脚本

以下是修改后的脚本,实现批量读取FlowFile并生成单条ClickHouse批量插入查询:

import groovy.json.JsonSlurper
import java.nio.charset.StandardCharsets
import org.apache.commons.io.IOUtils

// 转换MongoDB日期格式为ClickHouse兼容的DateTime格式
String toClickhouseDateFormat(String mongoDate) {
    if (!mongoDate) return null

    def parsedDate = null
    def formats = ["yyyy-MM-dd'T'HH:mm:ss.SSS'Z'", "yyyy-MM-dd'T'HH:mm:ss'Z'"]

    formats.each { format ->
        try {
            parsedDate = Date.parse(format, mongoDate)
            return
        } catch (Exception e) {
            // 忽略格式不匹配的异常,尝试下一种格式
        }
    }

    return parsedDate != null ? parsedDate.format("yyyy-MM-dd HH:mm:ss") : null
}

// SQL值转义处理
String sanitizeValue(Object value) {
    if (value == null) return "NULL"
    if (value instanceof String) {
        def strValue = value.trim()
        if (strValue.isEmpty()) return "NULL"
        return "'${strValue.replace("'", "''")}'"
    }
    // 数字、布尔等类型直接返回字符串形式
    return String.valueOf(value)
}

// 批量处理逻辑
def batchSize = 100 // 每次批量处理的FlowFile数量,可根据实际调整
def flowFiles = session.get(batchSize)
if (!flowFiles || flowFiles.isEmpty()) return

def valuesList = []
def failedFlowFiles = []

flowFiles.each { flowFile ->
    try {
        // 读取单个FlowFile的JSON内容
        def jsonContent = IOUtils.toString(session.read(flowFile), StandardCharsets.UTF_8)
        def jsonObject = new JsonSlurper().parseText(jsonContent)

        // 提取并处理字段
        def idValue = sanitizeValue(jsonObject?._id?.'$oid')
        def firstName = sanitizeValue(jsonObject?.first_name)
        def lastName = sanitizeValue(jsonObject?.last_name)
        // 如果有日期字段,示例:def createTime = sanitizeValue(toClickhouseDateFormat(jsonObject?.created_at?.'$date'))

        // 生成单条VALUES元组
        def valueTuple = "(${idValue}, ${firstName}, ${lastName})"
        valuesList.add(valueTuple)

        // 标记当前FlowFile为处理成功,后续会移除
        session.remove(flowFile)
    } catch (Exception e) {
        // 记录异常并将失败的FlowFile加入失败列表
        log.error("处理FlowFile ${flowFile.id}失败: ${e.message}", e)
        failedFlowFiles.add(flowFile)
    }
}

// 如果有成功处理的数据,生成批量插入语句
if (!valuesList.isEmpty()) {
    def insertQuery = """INSERT INTO db.table_name (id, first_name, last_name)
VALUES ${valuesList.join(',\n')};"""

    // 创建新的FlowFile保存批量插入语句
    def batchFlowFile = session.create()
    batchFlowFile = session.write(batchFlowFile, { outputStream ->
        outputStream.write(insertQuery.getBytes(StandardCharsets.UTF_8))
    } as OutputStreamCallback)

    // 转移批量生成的FlowFile到成功关系
    session.transfer(batchFlowFile, REL_SUCCESS)
}

// 转移处理失败的FlowFile到失败关系
if (!failedFlowFiles.isEmpty()) {
    session.transfer(failedFlowFiles, REL_FAILURE)
}

关键修改说明:

  • 批量获取FlowFile:使用session.get(batchSize)一次性获取多个FlowFile,batchSize可根据业务场景调整
  • 批量VALUES生成:遍历每个FlowFile解析JSON,将每条记录转换为VALUES元组,最后用逗号拼接成批量插入格式
  • 异常隔离:单个FlowFile解析失败时,不会中断整个批量处理流程,失败的FlowFile会被单独转移到REL_FAILURE关系
  • 优化值处理:扩展sanitizeValue函数,支持非字符串类型(数字、布尔等),直接返回符合SQL规范的格式
  • 资源清理:成功处理的FlowFile会被移除,避免重复处理

内容的提问来源于stack exchange,提问作者Ali Tariq

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 00:58:23