1. 需求说明
生产环境中有些数据需要在抽取的时候指定对某个字段进行过滤,判断等等。以将本地文件抽取到HDFS为例,当前我们需要导入的数据有2条,如下:

上面的数据中有uname字段,我们希望增加一个新的字段sex,该字段的值判断如果uname是wangwu,则sex字段的值就为female,否则为male,效果如下:

实现上面的效果需要2步:
- 编写过滤器代码。
- 将过滤器代码写到datax.json中。
2. 编写过滤器代码
-
导入datax的依赖(这里主要是因为要写日志,另一个是打包功能的配置,根据自己需要来添加依赖)
1<dependencies> 2 <dependency> 3 <groupId>org.slf4j</groupId> 4 <artifactId>slf4j-simple</artifactId> 5 <version>1.7.12</version> 6 </dependency> 7 </dependencies> 8 <build> 9 <finalName>gtmc-datax-utils-${project.version}</finalName> 10 <plugins> 11 <plugin> 12 <groupId>org.apache.maven.plugins</groupId> 13 <artifactId>maven-compiler-plugin</artifactId> 14 <version>3.7.0</version> 15 <configuration> 16 <source>1.7</source> 17 <target>1.7</target> 18 <encoding>UTF-8</encoding> 19 </configuration> 20 </plugin> 21 </plugins> 22 </build> -
编写过滤代码(注意是静态代码块,replaceQuotationMark方法为过滤方法,传递进一个字符串进行替换)
1public class StringUtils { 2 3 4 5 private static Logger LOG = LoggerFactory.getLogger(StringUtils.class); 6 7 /*** 8 * 替换双引号 9 * @param content 10 * @return 11 * @throws Exception 12 */ 13 public static String replaceQuotationMark(String content) throws Exception { 14 15 16 if (null == content || "".equals(content)) { 17 18 19 return ""; 20 } 21 try { 22 23 24 LOG.info("===" + content); 25 return content.equals("wangwu") ? "female" : "male"; 26 } catch (Exception e) { 27 28 29 LOG.error("替换双引号,失败原因:" + e.getMessage()); 30 return ""; 31 } 32 } 33} -
编译打包,将安装包丢到 datax/lib 目录下。
3. 应用过滤器
将插件应用到 datax.json文件中,过滤器的插件配置位置如下:
1{ 2 3 4 "job": { 5 6 7 "setting": { 8 9 ……}, 10 "content": [ 11 { 12 13 14 "reader": { 15 16 ……}, 17 "writer": { 18 19 ……}, 20 ## 过滤器位置 21 "transformer": [ 22 { 23 24 25 "name": "dx_substr", 26 "parameter": { 27 28 29 "columnIndex": 5, 30 "paras": [ "1", "3" ] 31 } 32 } 33 ] 34 } 35 ] 36 } 37}
我们的应用代码如下:
1"transformer": [ 2 { 3 4 5 ## 这里是固定的 6 "name": "dx_groovy", 7 "parameter": { 8 9 10 ## 这里定义数据是怎么跟我们的函数关联的 11 "code": "record.setColumn(4, new StringColumn(StringUtils.replaceQuotationMark(record.getColumn(4).asString())));\nreturn record;", 12 "extraPackage": [ 13 ## 导入函数所在的类,注意如下字符串有双引号,类的全路径后面有分号 14 "import com.gtmc.datax.utils.StringUtils;" 15 ] 16 } 17 } 18]