【问题标题】:Apache NiFi: Processing multiple csv's using the ExecuteScript ProcessorApache NiFi:使用 ExecuteScript 处理器处理多个 csv
【发布时间】:2020-03-01 14:09:05
【问题描述】:

我有一个 70 列的 csv。第 60 列包含一个值,该值决定记录是valid 还是invalid。如果第 60 列有 0、1、6 或 7,则为 valid。如果它包含任何其他值,那么它的invalid

我意识到这个功能不可能完全依赖于改变 Apache NiFi 中处理器的属性。因此,我决定使用executeScript processor 并将此 python 代码添加为文本正文。

import csv

valid =0
invalid =0
total =0
file2 = open("invalid.csv","w")
file1 = open("valid.csv","w")

with  open('/Users/himsaragallage/Desktop/redder/Regexo_2019101812750.dat.csv') as f:
    r = csv.reader(f)
    for row in f:
        # print row[1]
        total +=1

        if row[59] == "0" or row[59] == "1" or row[59] == "6" or row[59] == "7":
            valid +=1
            file1.write(row)
        else:
            invalid += 1
            file2.write(row)
file1.close()
file2.close()
print("Total : " + str(total))
print("Valid : " + str(valid))
print("Invalid : " + str(invalid))

我不知道如何在 executeScript 处理器中使用会话和代码,如this question 所示。所以我只是写了一个简单的python代码,并将有效和无效数据定向到不同的文件。我使用的这种方法有许多限制

  1. 我希望能够动态处理具有不同文件名的 csv。
  2. 发送无效数据的 csv 文件名也必须与输入 csv 文件名相同。
  3. 我的redder 文件夹中大约有 20 个 csv。必须一次性处理所有这些。

希望您能建议我执行以下操作的方法。随时通过编辑我使用的python代码甚至完全使用不同的处理器集并完全排除ExecuteScript Processer的使用来为我提供解决方案@

【问题讨论】:

  • 您可以查看 QueryRecord 处理器,而不是在 Jython 脚本中执行此操作。使用该处理器,您将能够简单地编写一个新关系,即“select * from FLOWFILE where column60 in (0,1,6,7)”
  • @Pushkr 您能否清楚在 QueryRecord 处理器中哪些属性/配置必须更改为“select * from FLOWFILE where column60 in (0,1,6,7) ”。

标签: python csv apache-nifi data-cleaning


【解决方案1】:

这是how to use QueryRecord processor上的完整分步说明

基本上,您需要设置突出显示的属性

【讨论】:

  • 这对我不起作用。您是否更改了CSVReaderCSVRecordSetWriter 中的任何配置?此外,我的 csv 中没有标题,因此我更改了 CSVRecordSetWriter 中的属性,以不将第一行作为标题。
  • 您是否在读取器和写入器中都指定了 csv 模式?你遇到了什么错误?
  • 它说标题包含重复的名称。但我已经确定我已将header 分配为fail
  • 如果您没有 CSV 标头,@HimsaraGallege QueryRecord 将不起作用。因为列名用作查询的一部分,例如col60。你有机会包含标题吗?否则我们必须考虑另一种解决方案。
  • @Upvote 我包含了标题并重试,这就是我遇到重复名称错误的时候。
【解决方案2】:

您希望根据一列中的值路由记录。在 NiFi 中有多种方法可以实现这一点。我可以想到以下几点:

我将向您展示如何使用PartitionRecord 处理器解决您的问题。由于您没有提供任何示例数据,我创建了一个示例用例。我想将欧洲的城市与其他地方的城市区分开来。给出以下数据:

id,city,country
1,Berlin,Germany
2,Paris,France
3,New York,USA
4,Frankfurt,Germany

流程:

生成流文件:

分区记录:

CSVReader 应设置为推断架构,CSVRecordSetWriter 应设置为继承架构。 PartitionRecord 将按国家/地区对记录进行分组,并将它们与具有国家/地区值的属性country 一起传递。您将看到以下记录组:

id,city,country
1,Berlin,Germany
4,Frankfurt,Germany

id,city,country
2,Paris,France

id,city,country
3,New York,USA

每个组都是一个流文件,并具有国家属性,您将使用该属性来路由组。

RouteOn属性:

来自欧洲的所有国家都将被路由到 is_europe 关系。现在您可以将相同的策略应用于您的用例。

【讨论】:

    猜你喜欢
    • 2020-03-10
    • 1970-01-01
    • 1970-01-01
    • 2016-08-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-01-11
    相关资源
    最近更新 更多