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

求Streamsets Jython Evaluator代码:检查记录所有字段含特定字符串并设属性

解决方案:递归遍历所有字段检查特定字符串

这里有一个实用的Jython实现,能帮你检查记录中**所有字段(包括嵌套的记录、列表)**是否包含目标字符串,并根据结果设置对应的头部属性:

# 替换成你需要检查的特定字符串
TARGET_STRING = "your_target_value"

def contains_target(value):
    # 处理列表类型的字段(比如数组)
    if isinstance(value, list):
        for item in value:
            if contains_target(item):
                return True
    # 处理嵌套的记录(字典结构)
    elif isinstance(value, dict):
        for field_value in value.values():
            if contains_target(field_value):
                return True
    # 处理字符串类型的字段值
    elif isinstance(value, str):
        if TARGET_STRING in value:
            return True
    # 非字符串类型(数字、布尔等)直接跳过
    return False

for record in records:
    try:
        # 获取整个记录的根节点数据
        full_record = record.getValue('/')
        if contains_target(full_record):
            record.attributes["DATA"] = "BAD"
        else:
            record.attributes["DATA"] = "GOOD"
        # 写入正常输出流
        sdc.output.write(record)
    except Exception as e:
        # 处理异常情况(比如记录结构异常)
        record.attributes["DATA"] = "ERROR"
        sdc.error.write(record, f"检查字段时出错: {str(e)}")

代码说明

  1. 递归遍历函数:contains_target会逐层遍历记录的所有内容,不管是嵌套的子记录还是数组列表,都能覆盖到;
  2. 获取完整记录:用record.getValue('/')替代record.values['/*'],这个方法是StreamSets官方推荐的获取节点数据的方式,/代表记录的根节点,能拿到完整的记录结构;
  3. 异常处理:添加了try-except块,避免因为记录格式异常导致整个管道中断,异常时会把标记设为ERROR并写入错误流;
  4. 灵活适配:如果需要忽略大小写,可以把字符串检查改成if TARGET_STRING.lower() in value.lower()。

为什么record.values['/*']不行?

record.values本质是一个Python字典,并不支持/*这种StreamSets的路径表达式语法。而record.getValue()是专门用来解析路径表达式的API,所以用它来获取根节点数据才是正确的做法。

内容的提问来源于stack exchange,提问作者john carlo estipona

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 13:52:29