Apache NiFi Validate CSV处理器可选字段验证报错及解决方案咨询
解决Apache NiFi Validate CSV必填/可选字段验证问题
一、Validate CSV处理器的正确Schema语法
NiFi的Validate CSV处理器Schema语法不支持单独的Optional()函数,你需要通过以下方式定义必填/可选字段:
- 必填字段:使用带非空约束的类型,比如
StrNotNullOrEmpty()(字符串非空)、ParseDate("MM/dd/yyyy")(日期必须有效且存在) - 可选字段:
- 允许字段存在但为空值:使用带问号的可空类型,比如
StrNotNullOrEmpty?()、ParseDate?("MM/dd/yyyy") - 允许字段完全缺失:使用
Str()(允许字段缺失或为空字符串),或者Schema中不包含该字段
- 允许字段存在但为空值:使用带问号的可空类型,比如
针对你的12个字段需求(第5、10、11位为可选),修正后的Schema示例:
StrNotNullOrEmpty(), ParseDate("MM/dd/yyyy"), StrNotNullOrEmpty(), StrNotNullOrEmpty(), Str(), StrNotNullOrEmpty(), StrNotNullOrEmpty(), StrNotNullOrEmpty(), StrNotNullOrEmpty(), Str(), Str(), StrNotNullOrEmpty()
如果需要可选字段必须存在但允许为空,用可空类型写法:
StrNotNullOrEmpty(), ParseDate("MM/dd/yyyy"), StrNotNullOrEmpty(), StrNotNullOrEmpty(), StrNotNullOrEmpty?(), StrNotNullOrEmpty(), StrNotNullOrEmpty(), StrNotNullOrEmpty(), StrNotNullOrEmpty(), StrNotNullOrEmpty?(), StrNotNullOrEmpty?(), StrNotNullOrEmpty()
二、Groovy验证脚本示例
如果内置Schema满足不了复杂逻辑,用ExecuteScript处理器运行以下Groovy脚本:
import org.apache.commons.csv.CSVFormat import org.apache.commons.csv.CSVParser import org.apache.commons.csv.CSVRecord def flowFile = session.get() if (!flowFile) return // 定义字段验证规则:索引对应验证逻辑 def fieldRules = [ 0: { v -> v != null && v.trim() != "" }, // 必填非空 1: { v -> if (!v) return true // 可选字段允许为空 try { Date.parse("MM/dd/yyyy", v); return true } catch (Exception e) { return false } }, 4: { v -> true }, // 完全可选无验证 9: { v -> true }, 10: { v -> true }, 2: { v -> v != null && v.trim() != "" }, 3: { v -> v != null && v.trim() != "" }, 5: { v -> v != null && v.trim() != "" }, 6: { v -> v != null && v.trim() != "" }, 7: { v -> v != null && v.trim() != "" }, 8: { v -> v != null && v.trim() != "" }, 11: { v -> v != null && v.trim() != "" } ] session.read(flowFile, { inputStream -> def reader = new InputStreamReader(inputStream, "UTF-8") def parser = CSVParser.parse(reader, CSVFormat.DEFAULT) // 有表头加.withHeader() def valid = true for (CSVRecord record : parser) { fieldRules.each { idx, rule -> def value = record.get(idx) if (!rule(value)) { valid = false log.error("字段索引${idx}验证失败,值: ${value}") } } if (!valid) break } flowFile = session.putAttribute(flowFile, "csv.valid", valid ? "true" : "false") session.transfer(flowFile, valid ? REL_SUCCESS : REL_FAILURE) } as InputStreamCallback)
三、Python验证脚本示例
同样用ExecuteScript处理器,选择Python运行以下脚本:
import csv from java.io import InputStreamReader, BufferedReader from org.apache.nifi.processor.io import InputStreamCallback class CSVValidator(InputStreamCallback): def __init__(self, session): self.session = session def process(self, inputStream): reader = BufferedReader(InputStreamReader(inputStream, "UTF-8")) csv_reader = csv.reader(reader) # 字段验证规则 def validate_date(v): if not v: return True try: from datetime import datetime datetime.strptime(v, "%m/%d/%Y") return True except: return False field_rules = [ lambda v: v is not None and v.strip() != "", # 索引0必填 validate_date, # 索引1可选日期 lambda v: v is not None and v.strip() != "", # 索引2必填 lambda v: v is not None and v.strip() != "", # 索引3必填 lambda v: True, # 索引4可选 lambda v: v is not None and v.strip() != "", # 索引5必填 lambda v: v is not None and v.strip() != "", # 索引6必填 lambda v: v is not None and v.strip() != "", # 索引7必填 lambda v: v is not None and v.strip() != "", # 索引8必填 lambda v: True, # 索引9可选 lambda v: True, # 索引10可选 lambda v: v is not None and v.strip() != "" # 索引11必填 ] valid = True for row_num, row in enumerate(csv_reader): for idx, rule in enumerate(field_rules): value = row[idx] if idx < len(row) else None if not rule(value): valid = False log.error(f"第{row_num+1}行,字段索引{idx}验证失败,值: {value}") break if not valid: break flowFile = self.session.get() if valid: flowFile = self.session.putAttribute(flowFile, "csv.valid", "true") self.session.transfer(flowFile, REL_SUCCESS) else: flowFile = self.session.putAttribute(flowFile, "csv.valid", "false") self.session.transfer(flowFile, REL_FAILURE) flowFile = session.get() if flowFile is not None: session.read(flowFile, CSVValidator(session))
内容的提问来源于stack exchange,提问作者Aiden Martin
相关产品推荐
相关产品推荐

