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

Apache NiFi Validate CSV处理器可选字段验证报错及解决方案咨询

解决Apache NiFi Validate CSV必填/可选字段验证问题

一、Validate CSV处理器的正确Schema语法

NiFi的Validate CSV处理器Schema语法不支持单独的Optional()函数,你需要通过以下方式定义必填/可选字段:

  • 必填字段:使用带非空约束的类型,比如StrNotNullOrEmpty()(字符串非空)、ParseDate("MM/dd/yyyy")(日期必须有效且存在)
  • 可选字段:
    1. 允许字段存在但为空值:使用带问号的可空类型,比如StrNotNullOrEmpty?()、ParseDate?("MM/dd/yyyy")
    2. 允许字段完全缺失:使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 19:23:19