求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)}")
代码说明
- 递归遍历函数:
contains_target会逐层遍历记录的所有内容,不管是嵌套的子记录还是数组列表,都能覆盖到; - 获取完整记录:用
record.getValue('/')替代record.values['/*'],这个方法是StreamSets官方推荐的获取节点数据的方式,/代表记录的根节点,能拿到完整的记录结构; - 异常处理:添加了try-except块,避免因为记录格式异常导致整个管道中断,异常时会把标记设为
ERROR并写入错误流; - 灵活适配:如果需要忽略大小写,可以把字符串检查改成
if TARGET_STRING.lower() in value.lower()。
为什么record.values['/*']不行?
record.values本质是一个Python字典,并不支持/*这种StreamSets的路径表达式语法。而record.getValue()是专门用来解析路径表达式的API,所以用它来获取根节点数据才是正确的做法。
内容的提问来源于stack exchange,提问作者john carlo estipona
相关产品推荐
相关产品推荐

