尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Flume拦截器实战:自定义ETL、数据脱敏与动态路由标签

Flume拦截器实战:自定义ETL、数据脱敏与动态路由标签 Flume拦截器实战自定义ETL、数据脱敏与动态路由标签1. 引言Flume作为Hadoop生态系统中的日志采集工具其拦截器(Interceptor)机制提供了强大的数据处理能力。本文将通过实战案例介绍如何自定义拦截器实现ETL转换、数据脱敏和动态路由标签功能提升数据采集与处理的灵活性和安全性。2. Flume拦截器基础Flume拦截器位于Source和Channel之间能够在数据进入Channel前对事件(Event)进行预处理。一个拦截器需要实现Interceptor接口主要包括以下方法initialize()初始化方法intercept(Event event)拦截单个事件intercept(ListEvent events)拦截事件列表close()关闭资源拦截器可以链式调用多个拦截器按照配置顺序依次执行。3. 自定义ETL拦截器实战ETL拦截器主要用于数据清洗、转换和标准化处理。下面是一个简单的ETL拦截器实现public class ETLInterceptor implements Interceptor { private String delimiter; // 分隔符 public ETLInterceptor(String delimiter) { this.delimiter delimiter; } Override public void initialize() { // 初始化逻辑 } Override public Event intercept(Event event) { // 处理单个事件 String body new String(event.getBody()); String[] fields body.split(delimiter); // 数据转换逻辑示例将时间戳转换为日期格式 if (fields.length 0) { long timestamp Long.parseLong(fields[0]); String date new SimpleDateFormat(yyyy-MM-dd HH:mm:ss).format(new Date(timestamp)); fields[0] date; } // 构建新事件体 String newBody String.join(delimiter, fields); event.setBody(newBody.getBytes()); return event; } Override public ListEvent intercept(ListEvent events) { // 处理事件列表 ListEvent interceptedEvents new ArrayList(); for (Event event : events) { interceptedEvents.add(intercept(event)); } return interceptedEvents; } Override public void close() { // 关闭资源 } // 构建器模式 public static class Builder implements Interceptor.Builder { private String delimiter ,; Override public void configure(Context context) { delimiter context.getString(delimiter, ,); } Override public Interceptor build() { return new ETLInterceptor(delimiter); } } }配置示例# 在flume.conf中配置拦截器 agent.sources.r1.interceptors i1 agent.sources.r1.interceptors.i1.type com.example.ETLInterceptor$Builder agent.sources.r1.interceptors.i1.delimiter ,4. 数据脱敏拦截器实战数据脱敏拦截器用于保护敏感信息如身份证号、手机号、邮箱等。下面是一个脱敏拦截器的实现public class DataMaskInterceptor implements Interceptor { private Pattern idCardPattern; // 身份证号正则 private Pattern phonePattern; // 手机号正则 private Pattern emailPattern; // 邮箱正则 public DataMaskInterceptor() { // 初始化正则表达式 idCardPattern Pattern.compile((\\d{6})\\d{8}(\\d{4})); phonePattern Pattern.compile((\\d{3})\\d{4}(\\d{4})); emailPattern Pattern.compile((.{2})(.*)(.*)); } Override public void initialize() { // 初始化逻辑 } Override public Event intercept(Event event) { // 处理单个事件 String body new String(event.getBody()); // 身份证脱敏保留前6位和后4位中间用*代替 body idCardPattern.matcher(body).replaceAll($1********$2); // 手机号脱敏保留前3位和后4位中间用*代替 body phonePattern.matcher(body).replaceAll($1****$2); // 邮箱脱敏保留前2位和后的内容中间用*代替 body emailPattern.matcher(body).replaceAll($1*$3); event.setBody(body.getBytes()); return event; } Override public ListEvent intercept(ListEvent events) { // 处理事件列表 ListEvent interceptedEvents new ArrayList(); for (Event event : events) { interceptedEvents.add(intercept(event)); } return interceptedEvents; } Override public void close() { // 关闭资源 } // 构建器模式 public static class Builder implements Interceptor.Builder { Override public void configure(Context context) { // 配置参数 } Override public Interceptor build() { return new DataMaskInterceptor(); } } }配置示例# 在flume.conf中配置拦截器 agent.sources.r1.interceptors i1 agent.sources.r1.interceptors.i1.type com.example.DataMaskInterceptor$Builder5. 动态路由标签拦截器实战动态路由标签拦截器用于根据数据内容添加标签实现数据分流。下面是一个动态路由标签拦截器的实现public class DynamicRouterInterceptor implements Interceptor { private MapString, String routeRules; // 路由规则 public DynamicRouterInterceptor(MapString, String routeRules) { this.routeRules routeRules; } Override public void initialize() { // 初始化逻辑 } Override public Event intercept(Event event) { // 处理单个事件 String body new String(event.getBody()); // 添加路由标签 for (Map.EntryString, String entry : routeRules.entrySet()) { if (body.contains(entry.getKey())) { event.getHeaders().put(router-tag, entry.getValue()); break; } } return event; } Override public ListEvent intercept(ListEvent events) { // 处理事件列表 ListEvent interceptedEvents new ArrayList(); for (Event event : events) { interceptedEvents.add(intercept(event)); } return interceptedEvents; } Override public void close() { // 关闭资源 } // 构建器模式 public static class Builder implements Interceptor.Builder { private MapString, String routeRules new HashMap(); Override public void configure(Context context) { // 从配置中加载路由规则 String rules context.getString(routeRules); String[] pairs rules.split(,); for (String pair : pairs) { String[] keyValue pair.split(); if (keyValue.length 2) { routeRules.put(keyValue[0], keyValue[1]); } } } Override public Interceptor build() { return new DynamicRouterInterceptor(routeRules); } } }配置示例# 在flume.conf中配置拦截器 agent.sources.r1.interceptors i1 agent.sources.r1.interceptors.i1.type com.example.DynamicRouterInterceptor$Builder agent.sources.r1.interceptors.i1.routeRules errorerror-log,warningwarning-log,infoinfo-logFlume拦截器执行流程Source采集数据拦截器链处理ETL拦截器数据脱敏拦截器动态路由标签拦截器Channel暂存数据Sink消费数据6. 最小示例与注意事项最小示例将以上三个拦截器组合使用创建一个完整的Flume配置文件。# 定义Source agent.sources.r1.type exec agent.sources.r1.command tail -F /var/log/application.log # 定义拦截器 agent.sources.r1.interceptors i1 i2 i3 agent.sources.r1.interceptors.i1.type com.example.ETLInterceptor$Builder agent.sources.r1.interceptors.i1.delimiter | agent.sources.r1.interceptors.i2.type com.example.DataMaskInterceptor$Builder agent.sources.r1.interceptors.i3.type com.example.DynamicRouterInterceptor$Builder agent.sources.r1.interceptors.i3.routeRules ERRORerror-channel,WARNwarning-channel # 定义Channel agent.channels.c1.type memory agent.channels.c1.capacity 1000 agent.channels.c1.transactionCapacity 100 agent.channels.error-channel.type memory agent.channels.error-channel.capacity 1000 agent.channels.error-channel.transactionCapacity 100 agent.channels.warning-channel.type memory agent.channels.warning-channel.capacity 1000 agent.channels.warning-channel.transactionCapacity 100 # 定义Sink agent.sinks.k1.type logger agent.sinks.k1.channel c1 agent.sinks.error-sink.type logger agent.sinks.error-sink.channel error-channel agent.sinks.warning-sink.type logger agent.sinks.warning-sink.channel warning-channel # 连接组件 agent.sources.r1.channels c1 error-channel warning-channel agent.sources.r1.selector.type multiplexing agent.sources.r1.selector.header router-tag agent.sources.r1.selector.error-channel error agent.sources.r1.selector.warning-channel warning注意事项拦截器顺序很重要ETL拦截器通常应该放在最前面数据脱敏处理会增加CPU开销对于大量数据需要考虑性能影响动态路由标签拦截器的规则不宜过多否则会影响处理效率自定义拦截器需要打包成jar文件并放到Flume的lib目录下复杂的拦截器逻辑应该考虑异常处理避免影响整个Flume流程
返回列表