如何修改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
相关产品推荐
相关产品推荐

