【问题标题】:Update csv value using executescript processor fails in apache-nifi在 apache-nifi 中使用 executescript 处理器更新 csv 值失败
【发布时间】:2020-03-10 00:34:03
【问题描述】:

我尝试从流文件中读取并使用 csv 中的默认值更新记录值。为此,我使用了 ExecuteScript 处理器,其中包含以下 python 代码。

import sys
import re
import traceback
from org.apache.commons.io import IOUtils
from org.apache.nifi.processor.io import StreamCallback
from org.python.core.util import StringUtil
from java.lang import Class
from java.io import BufferedReader
from java.io import InputStreamReader
from java.io import OutputStreamWriter

flowfile = session.get()
record = flowfile.getAttribute('record_type')

if record == '0':
    flowfile = session.putAttribute(flowfile,'record_type', 'NEW_USER')
    session.transfer(flowFile, REL_SUCCESS)
    session.commit()
elif record == '1':
    flowfile = session.putAttribute(flowfile,'record_type', 'OLD_USER')
    session.transfer(flowFile, REL_SUCCESS)
    session.commit()
else:
    flowfile = session.putAttribute(flowfile,'record_type', 'IGNORE')
    session.transfer(flowFile, REL_SUCCESS)
    session.commit()

writer.flush()
writer.close()
reader.close()

我的 csv 看起来像

id,record_type
1,0
2,1
3,2
4,0

结果应该是:

id,record_type
1,NEW_USER
2,OLD_USER
3,IGNORE
4,NEW_USER

我收到以下错误:

AttributeError : 'NoneType' 对象在中没有属性 'getAttribute' 第 13 行的脚本

上面写着record = flowfile.getAttribute('record_type')这是错误的..

我不知道如何解决这个问题,因为我不擅长 python

【问题讨论】:

  • ExecuteScript 处理整个文件(不是按记录)。 getAttribute 返回属性(如文件名)而不是内容。要更改内容,请使用 flowFile.write 函数。在inet 中搜索nifi python cookbook 并查看示例。
  • @daggett 感谢您的建议。但我仍然不明白如何获得一个值来比较。
  • 如果你不擅长python,使用记录处理可能会更好。检查 UpdateRecord 处理器。
  • @daggett 是的,我使用过UpdateRecord 处理器,但是如问题中所述,在一步替换多个值时遇到问题。
  • 您有基于记录的if,我认为在您的情况下可以使用 UpdateRecord。我可以展示如何为你的案例做 groovy 脚本..(我在 python 中也很糟糕)

标签: python csv apache-nifi


【解决方案1】:

这不是 python,但根据作者的评论可能是 groovy。

使用带有以下代码的 ExecuteGroovyScript 处理器:

def ff=session.get()
if(!ff)return

def map = [
    '0': 'NEW_USER',
    '1': 'OLD_USER',
]

ff.write{rawIn, rawOut->
    rawOut.withWriter("UTF-8"){w->
        rawIn.withReader("UTF-8"){r->
            int rowNum = 0
            //iterate lines from input stream and split each with coma
            r.splitEachLine( ',' ){row->
                if(rowNum>0){
                    //if not a header line then substitute value using map
                    row[1] = map[ row[1] ] ?: 'IGNORE'
                }
                //join and write row to output writer
                w << row.join(',') << '\n'
                rowNum++
            }
        }
    }
}

REL_SUCCESS << ff

【讨论】:

  • 我接受这个作为答案,因为它解决了我的问题。
  • 这是子问题,如果我们要比较string,我们需要做什么样的修改?
  • 在这个脚本中所有的字符串(除了 rowNum),即使它包含数字。 '0' 是一个字符串。
  • 我已经尝试使用定义为 use string from headers 的列,但它不像之前的示例那样工作。
  • 你的意思是你想从标题中按名称而不是按数字来评估列?
猜你喜欢
  • 1970-01-01
  • 2020-03-01
  • 2020-09-20
  • 1970-01-01
  • 2016-08-30
  • 1970-01-01
  • 2022-09-27
  • 1970-01-01
  • 2019-10-12
相关资源
最近更新 更多