【问题标题】:Nifi, CSV to Avro ValidateRecord - routing flowfile to failure because of double quote in field valueNifi,CSV 到 Avro ValidateRecord - 由于字段值中的双引号,将流文件路由到失败
【发布时间】:2019-12-24 14:26:27
【问题描述】:

我正在尝试使用 ValidateRecord 处理器将 CSV 验证为 Avro。 ValidateRecord 处理器的 Record Reader 属性设置为 CSVReader 控制器服务。此 CSVReader 控制器服务的引号字符设置为双引号 (")。

当我尝试验证流文件时,由于字段值中存在双引号,因此很少有流文件重定向到失败关系。

来自流文件内容的示例 csv 行:

"ICUA","01/22/2019","08:48:18",394846,"你删除了密钥吗?","YES---选择"接受回复" 然后继续删除","","","1"

我想使用 ReplaceText 但这会篡改字段的实际值。

如果有人可以提供处理这种情况的方法,那将非常有帮助。

谢谢!

【问题讨论】:

  • 要有一个有效的 cvs - 字符串中的每个双引号必须用额外的双引号转义。像这样:"YES---select ""Accept Response"" and continue with the remove"
  • 是的,但目前我们无法控制 csv 的生成。基本上它是由位于客户端站点的机器自动生成的。感谢您的评论!
  • 很公平,但您可能只剩下一个脚本解决方案或一个聪明的正则表达式,正如 daggett 建议的那样。否则,您使用的 ValidateRecord 可以正确识别行无效,但您希望它们有效。也许你可以使用 ValidateCsv 而不是 ValidateRecord,因为它有自己的 DSL 来定义规则
  • 我最终使用了 ExecuteScript 处理器。感谢您的建议。

标签: apache-nifi


【解决方案1】:

这不是一个功能齐全的解决方案。

可能正则表达式必须扩展。所以,这只是一个想法:

尝试使用带有以下参数的 ReplaceText:

search:       ([^,"])"([^,"])
replacement:  $1""$2

【讨论】:

  • 感谢您指明方向。正如 mattyb 建议的那样,显然似乎只有两个选项可以完成此任务 - ReplaceText 处理器和 ExecuteScript 进程。
【解决方案2】:

为了实现搜索和替换缺失的双引号,我使用了使用 Python 的 ExecuteScript 处理器,例如,

from org.apache.commons.io import IOUtils
from java.nio.charset import StandardCharsets
from org.apache.nifi.processor.io import StreamCallback
from org.apache.nifi.processors.script import ExecuteScript
from org.python.core.util.FileUtil import wrap
from io import StringIO
import re

# Define a subclass of StreamCallback for use in session.write()
class PyStreamCallback(StreamCallback):
    def __init__(self):
        pass

    def process(self, inputStream, outputStream):
        with wrap(inputStream) as f:
            lines = f.readlines()
            outer_new_value_list = []
            for csv_row in lines:
                field_value_list = csv_row.split('|')
                inner_new_value_list = []
                for field in field_value_list:
                    if field.count('"') > 2:
                        replaced_field = re.sub(r'(?!^|.$)["^]', '""', field)
                        inner_new_value_list.append(replaced_field)

                    else:
                        inner_new_value_list.append(field)
                row = '|'.join([str(elem) for elem in inner_new_value_list])
                outer_new_value_list.append(row)
            with wrap(outputStream, 'w') as filehandle:
                filehandle.writelines("%s" % line for line in outer_new_value_list)
# end class
flowFile = session.get()
if (flowFile != None):
    flowFile = session.write(flowFile, PyStreamCallback())
    session.transfer(flowFile, ExecuteScript.REL_SUCCESS)
# implicit return at the end

【讨论】:

    猜你喜欢
    • 2022-08-06
    • 1970-01-01
    • 1970-01-01
    • 2015-02-13
    • 1970-01-01
    • 2015-07-01
    • 1970-01-01
    • 1970-01-01
    • 2020-05-23
    相关资源
    最近更新 更多