Compare commits

...
Author SHA1 Message Date
Trafalgar 7ac827b597 update fastjson version
update fastjson version
2020-08-31 21:28:39 +08:00
Trafalgar f6afdefa9e Merge pull request #793 from heljoyLiu/master
gdbwriter: update column index prefix
2020-08-12 15:22:18 +08:00
Liu Jianping 6641150217 gdbwriter: update column index prefix 2020-08-12 15:18:33 +08:00
Trafalgar dee608b093 Update README.md
update DataX开源用户交流群 二维码
2020-08-12 12:41:59 +08:00
Trafalgar 6e2df5578e Add files via upload
add 交流群二维码图片
2020-08-12 12:27:20 +08:00
Trafalgar 066df492ca Update README.md
update  DataX开源用户交流群 信息
2020-08-12 11:59:15 +08:00
Trafalgar a9b77e4c98 Merge pull request #149 from springCat/master
官方的maven库tablestore-streamclient没有1.0.0-SNAPSHOT,只有1.0.0版本
2020-08-05 17:18:22 +08:00
Trafalgar 7666c5e8df Merge pull request #428 from zhouzf05/master
Fix otsstreamreader compile failure.
2020-08-05 17:12:57 +08:00
Trafalgar 7d239d5d78 Merge pull request #771 from lazlaz/master
修复rdbms插件,无法加载plugin.json下配置驱动,导致如虚谷、达梦等数据库连接报错问题
2020-08-05 17:11:06 +08:00
Trafalgar 024e7f41c1 Merge pull request #587 from BoYiZhang/master
fix bug  DataXException is not work .
2020-08-05 17:04:58 +08:00
Trafalgar 5419d044fe Merge pull request #486 from c-agam/patch-1
Update osswriter.md
2020-08-05 17:04:19 +08:00
Trafalgar 8583cb3e45 Merge pull request #229 from csurong/csurong_dev
增加odpsReader官方示例中的必须参数
2020-08-05 17:02:38 +08:00
Trafalgar f0c9e588ac Merge pull request #447 from lzrlizhirong/master
[fix]throw dataXException
2020-08-05 17:01:45 +08:00
Trafalgar 78f319ec4c Merge pull request #504 from woaer/patch-1
Update userGuid.md
2020-08-05 17:00:21 +08:00
Trafalgar 82f86a3ff1 Merge pull request #626 from sqdf1990/modify_dataxPluginDev_md
修改dataxPluginDev.md文件中的文字错误
2020-08-05 16:59:35 +08:00
Trafalgar cd6bbb6706 Merge pull request #703 from Galthen/patch-1
更新postgresqlwriter.md文件中关于"column": ["*"]中的*未显示的问题
2020-08-05 16:58:54 +08:00
Trafalgar 9fa49565ab Merge pull request #693 from RiverLiu/master
修复ClickhouseWriter缺少包的问题 #661 #676
2020-08-05 16:58:13 +08:00
Trafalgar 7c464ebc28 Merge pull request #692 from Willian-Zhang/patch-1
Update mysqlwriter.md
2020-08-05 16:57:35 +08:00
Trafalgar 7c1dca3ee4 Merge pull request #674 from duyong6380/master
为ftp输出文件增加后缀名注释
2020-08-05 16:56:56 +08:00
Trafalgar acb7e5cc19 Merge pull request #733 from tofuHero/patch-2
Update txtfilewriter.md
2020-08-05 16:55:46 +08:00
Trafalgar a375618ef9 Merge pull request #762 from randomGear/patch-1
Update oraclereader.md
2020-08-05 16:55:09 +08:00
Trafalgar 937abd5215 Merge pull request #754 from XuDaojie/dev-xdj
补充mongodbreader文档,修正部分文档错误。
2020-08-05 16:54:46 +08:00
Trafalgar 561d9002b7 Merge pull request #777 from coderxiao/patch-1
FQA->FAQ
2020-08-05 16:53:38 +08:00
coderxiao 4150aab4a9 FQA->FAQ
FQA->FAQ
2020-07-28 15:15:17 +08:00
liuzhai 04d0937779 修复rdbmswriter插件,无法plugin.json下驱动,导致数据库连接报错问题 2020-07-21 15:03:40 +08:00
liuzhai d763fa33c6 修复rdbmsreader插件,无法plugin.json下驱动,导致数据库连接报错问题 2020-07-21 15:03:26 +08:00
randomGear 1a2ca67bd7 Update oraclereader.md
在DEMO中的querysql在结尾不应该添加 `;`,不然会导致ORA-00911: 无效字符
2020-07-08 15:01:14 +08:00
XuDaojie 21f2539ce6 补充mongodbreader文档,修正部分文档错误。 2020-07-03 14:47:28 +08:00
tofuHero 66de85112d Update txtfilewriter.md
3.3 类型转换处 表格格式 markdown 不识别,已转成 可识别的样子
2020-06-19 15:24:24 +08:00
chenyang f0db8135e7 更新postgresqlwriter.md文件中关于*未展示的问题
源文件中的内容为:如果要依次写入全部列,使用*表示, 例如: "column": ["*"]
但是在网站上查看说明时,显示为:如果要依次写入全部列,使用表示, 例如: "column": [""]

这个引起了使用者的误解,对*进行转义后可显示正常,更改为:
如果要依次写入全部列,使用\*表示, 例如: "column": ["\*"]
2020-06-02 15:10:59 +08:00
Trafalgar e2aed2e9e0 Merge pull request #694 from dataccs/patch-2
Update odpswriter.md
2020-05-26 05:17:24 +08:00
dataccs f75321f4a4 Update odpswriter.md 2020-05-26 01:51:08 +08:00
Liu Jian 5891b49f5d Fix missing simulator package #661 #676 2020-05-25 16:10:38 +08:00
Willian Z f09515c95c Update mysqlwriter.md 2020-05-25 13:43:47 +08:00
duyong6380 b8428efe3c Update ftpwriter.md 2020-05-13 17:55:16 +08:00
Trafalgar 3b3fa878db Merge pull request #635 from heljoyLiu/pulgin-gdb-update-to-set-property
gdbwrtier:  support set-property
2020-04-23 11:48:37 +08:00
Trafalgar 0631ad67af Merge pull request #639 from heljoyLiu/new-plugin-gdbreader
reader: add gdbreader for GDB
2020-04-23 11:36:47 +08:00
Trafalgar 74b776cdc0 Merge pull request #644 from tujiye/datax_clickhouse_writer
Add Clickhouse Writer
2020-04-20 15:00:47 +08:00
jiye.tjy d9f2f4aa0d Add Clickhouse Writer 2020-04-14 21:22:23 +08:00
Liu Jianping 05d1851d99 gdbreader: reader for Aliyun GDB 2020-04-13 10:29:33 +08:00
Liu Jianping e09ec84f45 gdbwriter: update readme doc 2020-04-10 16:03:38 +08:00
Liu Jianping 2484343ade gdbwriter: update to support set-property
1. 添加id/label/属性字段长度限制
2. 添加对reader列索引格式,增加'#{i}'支持,同时兼容原格式
2. 添加SET属性导入支持
2020-04-10 15:43:28 +08:00
dufeng3 d6b70be5ac dataxPluginDev.md中把插件写成穿件了 2020-03-27 17:09:13 +08:00
dufeng3 07022e1276 修改dataxPluginDev.md文件中的文字错误 2020-03-27 16:10:51 +08:00
Trafalgar 643b6e9c64 Merge pull request #593 from dataccs/patch-1
Update osswriter.md
2020-02-17 20:02:07 +08:00
dataccs 9af6dfaf1d Update osswriter.md
fix the configuration of the osswriter sample
2020-02-17 19:56:27 +08:00
Trafalgar e1b1354a90 Update wiki for current support 2020-02-12 17:01:10 +08:00
BoYiZhang 235d4d3378 fix bug DataXException is not work . 2020-02-06 00:31:15 +08:00
Trafalgar f76d5cd921 Merge pull request #467 from wanda1416/master
更改 pom 版本号,解决外网用户无法编译通过。
review ok
2019-12-11 11:22:39 +08:00
wanda1416 c4bf7775a2 Merge branch 'master' into master 2019-11-19 17:39:43 +08:00
woaer 30cc3d56a2 Update userGuid.md
✏️ 修复拼写错误
2019-11-14 10:30:24 +08:00
Trafalgar a301cf5c6c Merge pull request #499 from asdf2014/github_readme
Add TSDB Reader into README doc
2019-11-08 20:30:12 +08:00
asdf2014 c2effde235 Add TSDB Reader into README doc 2019-11-08 20:14:55 +08:00
Trafalgar 928300f3cd Merge pull request #495 from asdf2014/github_tsdb_reader
Add TSDB Reader
2019-11-08 20:03:20 +08:00
asdf2014 dd19fd4332 Add TSDB Reader 2019-11-08 16:52:39 +08:00
c-agam ad3e8d6332 Update osswriter.md
object配置项修正:
经测试,原文档中object配置项:/cdo/datax应该替换为cdo/datax才能文档中所述功能
2019-10-31 11:34:21 +08:00
wanghui de093d73c6 更改 pom 版本号,解决外网用户无法编译通过。
datax pom 增加 aliyun 仓库镜像
  odpsreader 和 writer pom 修改 odps-sdk-core 为 0.20.7, 0.19.3 外网无法使用
  otsstreamreader pom 修改 client 为 1.0.0 正式版
2019-09-29 18:14:01 +08:00
lizhirong c64dab42aa [fix]throw dataXException 2019-09-20 14:16:25 +08:00
zhaofeng.zhou de583090dd Fix otsstreamreader compile failure. 2019-09-02 17:32:21 +08:00
csurong 24cdfc37ea 增加odpsReader官方示例中的必须参数 2018-12-03 17:25:13 +08:00
springcat 3944752098 官方的maven库tablestore-streamclient没有1.0.0-SNAPSHOT,只有1.0.0版本 2018-07-27 11:40:07 +08:00
92 changed files with 5980 additions and 1008 deletions
+17 -4
View File
@@ -57,16 +57,16 @@ DataX目前已经有了比较全面的插件体系,主流的RDBMS数据库、N
| | HDFS | √ | √ |[](https://github.com/alibaba/DataX/blob/master/hdfsreader/doc/hdfsreader.md) 、[](https://github.com/alibaba/DataX/blob/master/hdfswriter/doc/hdfswriter.md)|
| | Elasticsearch | | √ |[](https://github.com/alibaba/DataX/blob/master/elasticsearchwriter/doc/elasticsearchwriter.md)|
| 时间序列数据库 | OpenTSDB | √ | |[](https://github.com/alibaba/DataX/blob/master/opentsdbreader/doc/opentsdbreader.md)|
| | TSDB | | √ |[](https://github.com/alibaba/DataX/blob/master/tsdbwriter/doc/tsdbhttpwriter.md)|
| | TSDB | | √ |[](https://github.com/alibaba/DataX/blob/master/tsdbreader/doc/tsdbreader.md) 、[](https://github.com/alibaba/DataX/blob/master/tsdbwriter/doc/tsdbhttpwriter.md)|
# 我要开发新的插件
请点击:[DataX插件开发宝典](https://github.com/alibaba/DataX/blob/master/dataxPluginDev.md)
# 项目成员
核心Contributions: 光戈、一斅、祁然、云时
核心Contributions: 言柏 、枕水、秋奇、青砾、一斅、云时
感谢天烬、巴真、静行对DataX做出的贡献。
感谢天烬、光戈、祁然、巴真、静行对DataX做出的贡献。
# License
@@ -108,7 +108,20 @@ This software is free to use under the Apache License [Apache license](https://g
8. 对高并发、高稳定可用性、高性能、大数据处理有过实际项目及产品经验者优先考虑;
9. 有大数据产品、云产品、中间件技术解决方案者优先考虑。
````
钉钉用户群:23169395
钉钉用户群:
- DataX开源用户交流群
- <img src="https://github.com/alibaba/DataX/blob/master/images/DataX%E5%BC%80%E6%BA%90%E7%94%A8%E6%88%B7%E4%BA%A4%E6%B5%81%E7%BE%A4.jpg" width="20%" height="20%">
- DataX开源用户交流群2
- <img src="https://github.com/alibaba/DataX/blob/master/images/DataX%E5%BC%80%E6%BA%90%E7%94%A8%E6%88%B7%E4%BA%A4%E6%B5%81%E7%BE%A42.jpg" width="20%" height="20%">
- DataX开源用户交流群3
- <img src="https://github.com/alibaba/DataX/blob/master/images/DataX%E5%BC%80%E6%BA%90%E7%94%A8%E6%88%B7%E4%BA%A4%E6%B5%81%E7%BE%A43.jpg" width="20%" height="20%">
- DataX开源用户交流群4
- <img src="https://github.com/alibaba/DataX/blob/master/images/DataX%E5%BC%80%E6%BA%90%E7%94%A8%E6%88%B7%E4%BA%A4%E6%B5%81%E7%BE%A44.jpg" width="20%" height="20%">
- DataX开源用户交流群5
- <img src="https://github.com/alibaba/DataX/blob/master/images/DataX%E5%BC%80%E6%BA%90%E7%94%A8%E6%88%B7%E4%BA%A4%E6%B5%81%E7%BE%A45.jpg" width="20%" height="20%">
+88
View File
@@ -0,0 +1,88 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<artifactId>datax-all</artifactId>
<groupId>com.alibaba.datax</groupId>
<version>0.0.1-SNAPSHOT</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>clickhousewriter</artifactId>
<name>clickhousewriter</name>
<packaging>jar</packaging>
<dependencies>
<dependency>
<groupId>ru.yandex.clickhouse</groupId>
<artifactId>clickhouse-jdbc</artifactId>
<version>0.2.4</version>
</dependency>
<dependency>
<groupId>com.alibaba.datax</groupId>
<artifactId>datax-core</artifactId>
<version>${datax-project-version}</version>
</dependency>
<dependency>
<groupId>com.alibaba.datax</groupId>
<artifactId>datax-common</artifactId>
<version>${datax-project-version}</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
</dependency>
<dependency>
<groupId>ch.qos.logback</groupId>
<artifactId>logback-classic</artifactId>
</dependency>
<dependency>
<groupId>com.alibaba.datax</groupId>
<artifactId>plugin-rdbms-util</artifactId>
<version>${datax-project-version}</version>
</dependency>
</dependencies>
<build>
<resources>
<resource>
<directory>src/main/java</directory>
<includes>
<include>**/*.properties</include>
</includes>
</resource>
</resources>
<plugins>
<!-- compiler plugin -->
<plugin>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<source>${jdk-version}</source>
<target>${jdk-version}</target>
<encoding>${project-sourceEncoding}</encoding>
</configuration>
</plugin>
<!-- assembly plugin -->
<plugin>
<artifactId>maven-assembly-plugin</artifactId>
<configuration>
<descriptors>
<descriptor>src/main/assembly/package.xml</descriptor>
</descriptors>
<finalName>datax</finalName>
</configuration>
<executions>
<execution>
<id>dwzip</id>
<phase>package</phase>
<goals>
<goal>single</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
+35
View File
@@ -0,0 +1,35 @@
<assembly
xmlns="http://maven.apache.org/plugins/maven-assembly-plugin/assembly/1.1.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/plugins/maven-assembly-plugin/assembly/1.1.0 http://maven.apache.org/xsd/assembly-1.1.0.xsd">
<id></id>
<formats>
<format>dir</format>
</formats>
<includeBaseDirectory>false</includeBaseDirectory>
<fileSets>
<fileSet>
<directory>src/main/resources</directory>
<includes>
<include>plugin.json</include>
<include>plugin_job_template.json</include>
</includes>
<outputDirectory>plugin/writer/clickhousewriter</outputDirectory>
</fileSet>
<fileSet>
<directory>target/</directory>
<includes>
<include>clickhousewriter-0.0.1-SNAPSHOT.jar</include>
</includes>
<outputDirectory>plugin/writer/clickhousewriter</outputDirectory>
</fileSet>
</fileSets>
<dependencySets>
<dependencySet>
<useProjectArtifact>false</useProjectArtifact>
<outputDirectory>plugin/writer/clickhousewriter/libs</outputDirectory>
<scope>runtime</scope>
</dependencySet>
</dependencySets>
</assembly>
@@ -0,0 +1,329 @@
package com.alibaba.datax.plugin.writer.clickhousewriter;
import com.alibaba.datax.common.element.Column;
import com.alibaba.datax.common.element.StringColumn;
import com.alibaba.datax.common.exception.CommonErrorCode;
import com.alibaba.datax.common.exception.DataXException;
import com.alibaba.datax.common.plugin.RecordReceiver;
import com.alibaba.datax.common.spi.Writer;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.rdbms.util.DBUtilErrorCode;
import com.alibaba.datax.plugin.rdbms.util.DataBaseType;
import com.alibaba.datax.plugin.rdbms.writer.CommonRdbmsWriter;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONArray;
import java.sql.Array;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.sql.Timestamp;
import java.sql.Types;
import java.util.List;
import java.util.regex.Pattern;
public class ClickhouseWriter extends Writer {
private static final DataBaseType DATABASE_TYPE = DataBaseType.ClickHouse;
public static class Job extends Writer.Job {
private Configuration originalConfig = null;
private CommonRdbmsWriter.Job commonRdbmsWriterMaster;
@Override
public void init() {
this.originalConfig = super.getPluginJobConf();
this.commonRdbmsWriterMaster = new CommonRdbmsWriter.Job(DATABASE_TYPE);
this.commonRdbmsWriterMaster.init(this.originalConfig);
}
@Override
public void prepare() {
this.commonRdbmsWriterMaster.prepare(this.originalConfig);
}
@Override
public List<Configuration> split(int mandatoryNumber) {
return this.commonRdbmsWriterMaster.split(this.originalConfig, mandatoryNumber);
}
@Override
public void post() {
this.commonRdbmsWriterMaster.post(this.originalConfig);
}
@Override
public void destroy() {
this.commonRdbmsWriterMaster.destroy(this.originalConfig);
}
}
public static class Task extends Writer.Task {
private Configuration writerSliceConfig;
private CommonRdbmsWriter.Task commonRdbmsWriterSlave;
@Override
public void init() {
this.writerSliceConfig = super.getPluginJobConf();
this.commonRdbmsWriterSlave = new CommonRdbmsWriter.Task(DATABASE_TYPE) {
@Override
protected PreparedStatement fillPreparedStatementColumnType(PreparedStatement preparedStatement, int columnIndex, int columnSqltype, Column column) throws SQLException {
try {
if (column.getRawData() == null) {
preparedStatement.setNull(columnIndex + 1, columnSqltype);
return preparedStatement;
}
java.util.Date utilDate;
switch (columnSqltype) {
case Types.CHAR:
case Types.NCHAR:
case Types.CLOB:
case Types.NCLOB:
case Types.VARCHAR:
case Types.LONGVARCHAR:
case Types.NVARCHAR:
case Types.LONGNVARCHAR:
preparedStatement.setString(columnIndex + 1, column
.asString());
break;
case Types.TINYINT:
case Types.SMALLINT:
case Types.INTEGER:
case Types.BIGINT:
case Types.DECIMAL:
case Types.FLOAT:
case Types.REAL:
case Types.DOUBLE:
String strValue = column.asString();
if (emptyAsNull && "".equals(strValue)) {
preparedStatement.setNull(columnIndex + 1, columnSqltype);
} else {
switch (columnSqltype) {
case Types.TINYINT:
case Types.SMALLINT:
case Types.INTEGER:
preparedStatement.setInt(columnIndex + 1, column.asBigInteger().intValue());
break;
case Types.BIGINT:
preparedStatement.setLong(columnIndex + 1, column.asLong());
break;
case Types.DECIMAL:
preparedStatement.setBigDecimal(columnIndex + 1, column.asBigDecimal());
break;
case Types.REAL:
case Types.FLOAT:
preparedStatement.setFloat(columnIndex + 1, column.asDouble().floatValue());
break;
case Types.DOUBLE:
preparedStatement.setDouble(columnIndex + 1, column.asDouble());
break;
}
}
break;
case Types.DATE:
if (this.resultSetMetaData.getRight().get(columnIndex)
.equalsIgnoreCase("year")) {
if (column.asBigInteger() == null) {
preparedStatement.setString(columnIndex + 1, null);
} else {
preparedStatement.setInt(columnIndex + 1, column.asBigInteger().intValue());
}
} else {
java.sql.Date sqlDate = null;
try {
utilDate = column.asDate();
} catch (DataXException e) {
throw new SQLException(String.format(
"Date 类型转换错误:[%s]", column));
}
if (null != utilDate) {
sqlDate = new java.sql.Date(utilDate.getTime());
}
preparedStatement.setDate(columnIndex + 1, sqlDate);
}
break;
case Types.TIME:
java.sql.Time sqlTime = null;
try {
utilDate = column.asDate();
} catch (DataXException e) {
throw new SQLException(String.format(
"Date 类型转换错误:[%s]", column));
}
if (null != utilDate) {
sqlTime = new java.sql.Time(utilDate.getTime());
}
preparedStatement.setTime(columnIndex + 1, sqlTime);
break;
case Types.TIMESTAMP:
Timestamp sqlTimestamp = null;
if (column instanceof StringColumn && column.asString() != null) {
String timeStampStr = column.asString();
// JAVA TIMESTAMP 类型入参必须是 "2017-07-12 14:39:00.123566" 格式
String pattern = "^\\d+-\\d+-\\d+ \\d+:\\d+:\\d+.\\d+";
boolean isMatch = Pattern.matches(pattern, timeStampStr);
if (isMatch) {
sqlTimestamp = Timestamp.valueOf(timeStampStr);
preparedStatement.setTimestamp(columnIndex + 1, sqlTimestamp);
break;
}
}
try {
utilDate = column.asDate();
} catch (DataXException e) {
throw new SQLException(String.format(
"Date 类型转换错误:[%s]", column));
}
if (null != utilDate) {
sqlTimestamp = new Timestamp(
utilDate.getTime());
}
preparedStatement.setTimestamp(columnIndex + 1, sqlTimestamp);
break;
case Types.BINARY:
case Types.VARBINARY:
case Types.BLOB:
case Types.LONGVARBINARY:
preparedStatement.setBytes(columnIndex + 1, column
.asBytes());
break;
case Types.BOOLEAN:
preparedStatement.setInt(columnIndex + 1, column.asBigInteger().intValue());
break;
// warn: bit(1) -> Types.BIT 可使用setBoolean
// warn: bit(>1) -> Types.VARBINARY 可使用setBytes
case Types.BIT:
if (this.dataBaseType == DataBaseType.MySql) {
Boolean asBoolean = column.asBoolean();
if (asBoolean != null) {
preparedStatement.setBoolean(columnIndex + 1, asBoolean);
} else {
preparedStatement.setNull(columnIndex + 1, Types.BIT);
}
} else {
preparedStatement.setString(columnIndex + 1, column.asString());
}
break;
default:
boolean isHandled = fillPreparedStatementColumnType4CustomType(preparedStatement,
columnIndex, columnSqltype, column);
if (isHandled) {
break;
}
throw DataXException
.asDataXException(
DBUtilErrorCode.UNSUPPORTED_TYPE,
String.format(
"您的配置文件中的列配置信息有误. 因为DataX 不支持数据库写入这种字段类型. 字段名:[%s], 字段类型:[%d], 字段Java类型:[%s]. 请修改表中该字段的类型或者不同步该字段.",
this.resultSetMetaData.getLeft()
.get(columnIndex),
this.resultSetMetaData.getMiddle()
.get(columnIndex),
this.resultSetMetaData.getRight()
.get(columnIndex)));
}
return preparedStatement;
} catch (DataXException e) {
// fix类型转换或者溢出失败时,将具体哪一列打印出来
if (e.getErrorCode() == CommonErrorCode.CONVERT_NOT_SUPPORT ||
e.getErrorCode() == CommonErrorCode.CONVERT_OVER_FLOW) {
throw DataXException
.asDataXException(
e.getErrorCode(),
String.format(
"类型转化错误. 字段名:[%s], 字段类型:[%d], 字段Java类型:[%s]. 请修改表中该字段的类型或者不同步该字段.",
this.resultSetMetaData.getLeft()
.get(columnIndex),
this.resultSetMetaData.getMiddle()
.get(columnIndex),
this.resultSetMetaData.getRight()
.get(columnIndex)));
} else {
throw e;
}
}
}
private Object toJavaArray(Object val) {
if (null == val) {
return null;
} else if (val instanceof JSONArray) {
Object[] valArray = ((JSONArray) val).toArray();
for (int i = 0; i < valArray.length; i++) {
valArray[i] = this.toJavaArray(valArray[i]);
}
return valArray;
} else {
return val;
}
}
boolean fillPreparedStatementColumnType4CustomType(PreparedStatement ps,
int columnIndex, int columnSqltype,
Column column) throws SQLException {
switch (columnSqltype) {
case Types.OTHER:
if (this.resultSetMetaData.getRight().get(columnIndex).startsWith("Tuple")) {
throw DataXException
.asDataXException(ClickhouseWriterErrorCode.TUPLE_NOT_SUPPORTED_ERROR, ClickhouseWriterErrorCode.TUPLE_NOT_SUPPORTED_ERROR.getDescription());
} else {
ps.setString(columnIndex + 1, column.asString());
}
return true;
case Types.ARRAY:
Connection conn = ps.getConnection();
List<Object> values = JSON.parseArray(column.asString(), Object.class);
for (int i = 0; i < values.size(); i++) {
values.set(i, this.toJavaArray(values.get(i)));
}
Array array = conn.createArrayOf("String", values.toArray());
ps.setArray(columnIndex + 1, array);
return true;
default:
break;
}
return false;
}
};
this.commonRdbmsWriterSlave.init(this.writerSliceConfig);
}
@Override
public void prepare() {
this.commonRdbmsWriterSlave.prepare(this.writerSliceConfig);
}
@Override
public void startWrite(RecordReceiver recordReceiver) {
this.commonRdbmsWriterSlave.startWrite(recordReceiver, this.writerSliceConfig, super.getTaskPluginCollector());
}
@Override
public void post() {
this.commonRdbmsWriterSlave.post(this.writerSliceConfig);
}
@Override
public void destroy() {
this.commonRdbmsWriterSlave.destroy(this.writerSliceConfig);
}
}
}
@@ -0,0 +1,31 @@
package com.alibaba.datax.plugin.writer.clickhousewriter;
import com.alibaba.datax.common.spi.ErrorCode;
public enum ClickhouseWriterErrorCode implements ErrorCode {
TUPLE_NOT_SUPPORTED_ERROR("ClickhouseWriter-00", "不支持TUPLE类型导入."),
;
private final String code;
private final String description;
private ClickhouseWriterErrorCode(String code, String description) {
this.code = code;
this.description = description;
}
@Override
public String getCode() {
return this.code;
}
@Override
public String getDescription() {
return this.description;
}
@Override
public String toString() {
return String.format("Code:[%s], Description:[%s].", this.code, this.description);
}
}
+6
View File
@@ -0,0 +1,6 @@
{
"name": "clickhousewriter",
"class": "com.alibaba.datax.plugin.writer.clickhousewriter.ClickhouseWriter",
"description": "useScene: prod. mechanism: Jdbc connection using the database, execute insert sql.",
"developer": "jiye.tjy"
}
@@ -0,0 +1,21 @@
{
"name": "clickhousewriter",
"parameter": {
"username": "username",
"password": "password",
"column": ["col1", "col2", "col3"],
"connection": [
{
"jdbcUrl": "jdbc:clickhouse://<host>:<port>[/<database>]",
"table": ["table1", "table2"]
}
],
"preSql": [],
"postSql": [],
"batchSize": 65536,
"batchByteSize": 134217728,
"dryRun": false,
"writeMode": "insert"
}
}
@@ -427,7 +427,7 @@ public class JobContainer extends AbstractContainer {
Long channelLimitedByteSpeed = this.configuration
.getLong(CoreConstant.DATAX_CORE_TRANSPORT_CHANNEL_SPEED_BYTE);
if (channelLimitedByteSpeed == null || channelLimitedByteSpeed <= 0) {
DataXException.asDataXException(
throw DataXException.asDataXException(
FrameworkErrorCode.CONFIG_ERROR,
"在有总bps限速条件下,单个channel的bps值不能为空,也不能为非正数");
}
@@ -448,7 +448,7 @@ public class JobContainer extends AbstractContainer {
Long channelLimitedRecordSpeed = this.configuration.getLong(
CoreConstant.DATAX_CORE_TRANSPORT_CHANNEL_SPEED_RECORD);
if (channelLimitedRecordSpeed == null || channelLimitedRecordSpeed <= 0) {
DataXException.asDataXException(FrameworkErrorCode.CONFIG_ERROR,
throw DataXException.asDataXException(FrameworkErrorCode.CONFIG_ERROR,
"在有总tps限速条件下,单个channel的tps值不能为空,也不能为非正数");
}
+4 -4
View File
@@ -111,7 +111,7 @@ public class SomeReader extends Reader {
```
`Job`接口功能如下:
- `init`: Job对象初始化工作,测试可以通过`super.getPluginJobConf()`获取与本插件相关的配置。读插件获得配置中`reader`部分,写插件获得`writer`部分。
- `init`: Job对象初始化工作,此时可以通过`super.getPluginJobConf()`获取与本插件相关的配置。读插件获得配置中`reader`部分,写插件获得`writer`部分。
- `prepare`: 全局准备工作,比如odpswriter清空目标表。
- `split`: 拆分`Task`。参数`adviceNumber`框架建议的拆分数,一般是运行时所配置的并发度。值返回的是`Task`的配置列表。
- `post`: 全局的后置工作,比如mysqlwriter同步完影子表后的rename操作。
@@ -155,7 +155,7 @@ public class SomeReader extends Reader {
```
- `name`: 插件名称,大小写敏感。框架根据用户在配置文件中指定的名称来搜寻插件。 **十分重要**
- `class`: 入口类的全限定名称,框架通过反射穿件入口类的实例。**十分重要** 。
- `class`: 入口类的全限定名称,框架通过反射件入口类的实例。**十分重要** 。
- `description`: 描述信息。
- `developer`: 开发人员。
@@ -435,7 +435,7 @@ DataX的内部类型在实现上会选用不同的java类型:
#### 如何处理脏数据
在`Reader.Task`和`Writer.Task`中,过`AbstractTaskPlugin.getPluginCollector()`可以拿到一个`TaskPluginCollector`,它提供了一系列`collectDirtyRecord`的方法。当脏数据出现时,只需要调用合适的`collectDirtyRecord`方法,把被认为是脏数据的`Record`传入即可。
在`Reader.Task`和`Writer.Task`中,过`AbstractTaskPlugin.getTaskPluginCollector()`可以拿到一个`TaskPluginCollector`,它提供了一系列`collectDirtyRecord`的方法。当脏数据出现时,只需要调用合适的`collectDirtyRecord`方法,把被认为是脏数据的`Record`传入即可。
用户可以在任务的配置中指定脏数据限制条数或者百分比限制,当脏数据超出限制时,框架会结束同步任务,退出。插件需要保证脏数据都被收集到,其他工作交给框架就好。
@@ -468,4 +468,4 @@ DataX的内部类型在实现上会选用不同的java类型:
- 测试参数集(多组),系统参数(比如并发数),插件参数(比如batchSize)
- 不同参数下同步速度(Rec/s, MB/s),机器负载(load, cpu)等,对数据源压力(load, cpu, mem等)。
6. **约束限制**:是否存在其他的使用限制条件。
7. **FQA**:用户经常会遇到的问题。
7. **FAQ**:用户经常会遇到的问题。
+9
View File
@@ -63,6 +63,7 @@ FtpWriter实现了从DataX协议转为FTP文件功能,FTP文件本身是无结
"nullFormat": "null",
"dateFormat": "yyyy-MM-dd",
"fileFormat": "csv",
"suffix": ".csv",
"header": []
}
}
@@ -200,6 +201,14 @@ FtpWriter实现了从DataX协议转为FTP文件功能,FTP文件本身是无结
* 必选:否 <br />
* 默认值:text <br />
* **suffix**
* 描述:最后输出文件的后缀,当前支持 ".text"以及".csv"
* 必选:否 <br />
* 默认值:"" <br />
* **header**
+260
View File
@@ -0,0 +1,260 @@
# DataX GDBReader
## 1. 快速介绍
GDBReader插件实现读取GDB实例数据的功能,通过`Gremlin Client`连接远程GDB实例,按配置提供的`label`生成查询DSL,遍历点或边数据,包括属性数据,并将数据写入到Record中给到Writer使用。
## 2. 实现原理
GDBReader使用`Gremlin Client`连接GDB实例,按`label`分不同Task取点或边数据。
单个Task中按`label`遍历点或边的id,再切分范围分多次请求查询点或边和属性数据,最后将点或边数据根据配置转换成指定格式记录发送给下游写插件。
GDBReader按`label`切分多个Task并发,同一个`label`的数据批量异步获取来加快读取速度。如果配置读取的`label`列表为空,任务启动前会从GDB查询所有`label`再切分Task。
## 3. 功能说明
GDB中点和边不同,读取需要区分点和边点配置。
### 3.1 点配置样例
```
{
"job": {
"setting": {
"speed": {
"channel": 1
}
"errorLimit": {
"record": 1
}
},
"content": [
{
"reader": {
"name": "gdbreader",
"parameter": {
"host": "10.218.145.24",
"port": 8182,
"username": "***",
"password": "***",
"fetchBatchSize": 100,
"rangeSplitSize": 1000,
"labelType": "VERTEX",
"labels": ["label1", "label2"],
"column": [
{
"name": "id",
"type": "string",
"columnType": "primaryKey"
},
{
"name": "label",
"type": "string",
"columnType": "primaryLabel"
},
{
"name": "age",
"type": "int",
"columnType": "vertexProperty"
}
]
}
},
"writer": {
"name": "streamwriter",
"parameter": {
"print": true
}
}
}
]
}
}
```
### 3.2 边配置样例
```
{
"job": {
"setting": {
"speed": {
"channel": 1
},
"errorLimit": {
"record": 1
}
},
"content": [
{
"reader": {
"name": "gdbreader",
"parameter": {
"host": "10.218.145.24",
"port": 8182,
"username": "***",
"password": "***",
"fetchBatchSize": 100,
"rangeSplitSize": 1000,
"labelType": "EDGE",
"labels": ["label1", "label2"],
"column": [
{
"name": "id",
"type": "string",
"columnType": "primaryKey"
},
{
"name": "label",
"type": "string",
"columnType": "primaryLabel"
},
{
"name": "srcId",
"type": "string",
"columnType": "srcPrimaryKey"
},
{
"name": "srcLabel",
"type": "string",
"columnType": "srcPrimaryLabel"
},
{
"name": "dstId",
"type": "string",
"columnType": "srcPrimaryKey"
},
{
"name": "dstLabel",
"type": "string",
"columnType": "srcPrimaryLabel"
},
{
"name": "name",
"type": "string",
"columnType": "edgeProperty"
},
{
"name": "weight",
"type": "double",
"columnType": "edgeProperty"
}
]
}
},
"writer": {
"name": "streamwriter",
"parameter": {
"print": true
}
}
}
]
}
}
```
### 3.3 参数说明
* **host**
* 描述:GDB实例连接地址,对应'实例管理'->'基本信息'页面的网络地址
* 必选:是
* 默认值:无
* **port**
* 描述:GDB实例连接地址对应的端口
* 必选:是
* 默认值:8182
* **username**
* 描述:GDB实例账号名
* 必选:是
* 默认值:无
* **password**
* 描述:GDB实例账号名对应的密码
* 必选:是
* 默认值:无
* **fetchBatchSize**
* 描述:一次GDB请求读取点或边的数量,响应包含点或边以及属性
* 必选:是
* 默认值:100
* **rangeSplitSize**
* 描述:id遍历,一次遍历请求扫描的id个数
* 必选:是
* 默认值:10 \* fetchBatchSize
* **labels**
* 描述:标签数组,即需要导出的点或边标签,支持读取多个标签,用数组表示。如果留空([]),表示GDB中所有点或边标签
* 必选:是
* 默认值:无
* **labelType**
* 描述:数据标签类型,支持点、边两种枚举值
* VERTEX:表示点
* EDGE:表示边
* 必选:是
* 默认值:无
* **column**
* 描述:点或边字段映射关系配置
* 必选:是
* 默认值:无
* **column -> name**
* 描述:点或边映射关系的字段名,指定属性时表示读取的属性名,读取其他字段时会被忽略
* 必选:是
* 默认值:无
* **column -> type**
* 描述:点或边映射关系的字段类型
* id, label在GDB中都是string类型,配置非string类型时可能会转换失败
* 普通属性支持基础类型,包括int, long, float, double, boolean, string
* GDBReader尽量将读取到的数据转换成配置要求的类型,但转换失败会导致该条记录错误
* 必选:是
* 默认值:无
* **column -> columnType**
* 描述:GDB点或边数据到列数据的映射关系,支持以下枚举值:
* primaryKey 表示该字段是点或边的id
* primaryLabel 表示该字段是点或边的label
* srcPrimaryKey: 表示该字段是边关联的起点id,只在读取边时使用
* srcPrimaryLabel 表示该字段是边关联的起点label,只在读取边时使用
* dstPrimaryKey: 表示该字段是边关联的终点id,只在读取边时使用
* dstPrimaryLabel 表示该字段是边关联的终点label,只在读取边时使用
* vertexProperty: 表示该字段是点的属性,只在读取点时使用,应用到SET属性时只读取其中的一个属性值
* vertexJsonProperty: 表示该字段是点的属性集合,只在读取点时使用。属性集合使用JSON格式输出,包含所有的属性,不能与其他vertexProperty配置一起使用
* edgeProperty: 表示该字段是边的属性,只在读取边时使用
* edgeJsonProperty: 表示该字段是边的属性集合,只在读取边时使用。属性集合使用JSON格式输出,包含所有的属性,不能与其他edgeProperty配置一起使用
* 必选:是
* 默认值:无
* vertexJsonProperty格式示例,新增`c`字段区分SET属性,但是SET属性只包含单个属性值时会标记成普通属性
```
{"properties":[
{"k":"name","t","string","v":"Jack","c":"set"},
{"k":"name","t","string","v":"Luck","c":"set"},
{"k":"age","t","int","v":"20","c":"single"}
]}
```
* edgeJsonProperty格式示例,边不支持多值属性
```
{"properties":[
{"k":"created_at","t","long","v":"153498653"},
{"k":"weight","t","double","v":"3.14"}
]}
## 4 性能报告
(TODO)
## 5 使用约束
## 6 FAQ
+125
View File
@@ -0,0 +1,125 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<artifactId>datax-all</artifactId>
<groupId>com.alibaba.datax</groupId>
<version>0.0.1-SNAPSHOT</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>gdbreader</artifactId>
<groupId>com.alibaba.datax</groupId>
<version>0.0.1-SNAPSHOT</version>
<dependencies>
<dependency>
<groupId>com.alibaba.datax</groupId>
<artifactId>datax-common</artifactId>
<version>${datax-project-version}</version>
<exclusions>
<exclusion>
<artifactId>slf4j-log4j12</artifactId>
<groupId>org.slf4j</groupId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>com.alibaba.datax</groupId>
<artifactId>datax-core</artifactId>
<version>${datax-project-version}</version>
<scope>test</scope>
<exclusions>
<exclusion>
<artifactId>slf4j-log4j12</artifactId>
<groupId>org.slf4j</groupId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
</dependency>
<dependency>
<groupId>ch.qos.logback</groupId>
<artifactId>logback-classic</artifactId>
</dependency>
<dependency>
<groupId>org.apache.tinkerpop</groupId>
<artifactId>gremlin-driver</artifactId>
<version>3.4.1</version>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<version>1.18.8</version>
</dependency>
<dependency>
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter-api</artifactId>
<version>5.4.0</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter-engine</artifactId>
<version>5.4.0</version>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
<!-- compiler plugin -->
<plugin>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<source>1.6</source>
<target>1.6</target>
<encoding>${project-sourceEncoding}</encoding>
</configuration>
</plugin>
<!-- test case plugin -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-surefire-plugin</artifactId>
<version>2.22.0</version>
<configuration>
<includes>
<include>**/*Test*.class</include>
</includes>
</configuration>
</plugin>
<!-- assembly plugin -->
<plugin>
<artifactId>maven-assembly-plugin</artifactId>
<configuration>
<descriptors>
<descriptor>src/main/assembly/package.xml</descriptor>
</descriptors>
<finalName>datax</finalName>
</configuration>
<executions>
<execution>
<id>dwzip</id>
<phase>package</phase>
<goals>
<goal>single</goal>
</goals>
</execution>
</executions>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<source>8</source>
<target>8</target>
</configuration>
</plugin>
</plugins>
</build>
</project>
+35
View File
@@ -0,0 +1,35 @@
<assembly
xmlns="http://maven.apache.org/plugins/maven-assembly-plugin/assembly/1.1.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/plugins/maven-assembly-plugin/assembly/1.1.0 http://maven.apache.org/xsd/assembly-1.1.0.xsd">
<id></id>
<formats>
<format>dir</format>
</formats>
<includeBaseDirectory>false</includeBaseDirectory>
<fileSets>
<fileSet>
<directory>src/main/resources</directory>
<includes>
<include>plugin.json</include>
<include>plugin_job_template.json</include>
</includes>
<outputDirectory>plugin/reader/gdbreader</outputDirectory>
</fileSet>
<fileSet>
<directory>target/</directory>
<includes>
<include>gdbreader-0.0.1-SNAPSHOT.jar</include>
</includes>
<outputDirectory>plugin/reader/gdbreader</outputDirectory>
</fileSet>
</fileSets>
<dependencySets>
<dependencySet>
<useProjectArtifact>false</useProjectArtifact>
<outputDirectory>plugin/reader/gdbreader/libs</outputDirectory>
<scope>runtime</scope>
</dependencySet>
</dependencySets>
</assembly>
@@ -0,0 +1,231 @@
package com.alibaba.datax.plugin.reader.gdbreader;
import com.alibaba.datax.common.element.Record;
import com.alibaba.datax.common.exception.DataXException;
import com.alibaba.datax.common.plugin.RecordSender;
import com.alibaba.datax.common.spi.Reader;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.reader.gdbreader.mapping.DefaultGdbMapper;
import com.alibaba.datax.plugin.reader.gdbreader.mapping.MappingRule;
import com.alibaba.datax.plugin.reader.gdbreader.mapping.MappingRuleFactory;
import com.alibaba.datax.plugin.reader.gdbreader.model.GdbElement;
import com.alibaba.datax.plugin.reader.gdbreader.model.GdbGraph;
import com.alibaba.datax.plugin.reader.gdbreader.model.ScriptGdbGraph;
import com.alibaba.datax.plugin.reader.gdbreader.util.ConfigHelper;
import org.apache.tinkerpop.gremlin.driver.ResultSet;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.LinkedList;
import java.util.List;
public class GdbReader extends Reader {
private final static int DEFAULT_FETCH_BATCH_SIZE = 200;
private static GdbGraph graph;
private static Key.ExportType exportType;
/**
* Job 中的方法仅执行一次,Task 中方法会由框架启动多个 Task 线程并行执行。
* <p/>
* 整个 Reader 执行流程是:
* <pre>
* Job类init-->prepare-->split
*
* Task类init-->prepare-->startRead-->post-->destroy
* Task类init-->prepare-->startRead-->post-->destroy
*
* Job类post-->destroy
* </pre>
*/
public static class Job extends Reader.Job {
private static final Logger LOG = LoggerFactory.getLogger(Job.class);
private Configuration jobConfig = null;
@Override
public void init() {
this.jobConfig = super.getPluginJobConf();
/**
* 注意:此方法仅执行一次。
* 最佳实践:通常在这里对用户的配置进行校验:是否缺失必填项?有无错误值?有没有无关配置项?...
* 并给出清晰的报错/警告提示。校验通常建议采用静态工具类进行,以保证本类结构清晰。
*/
ConfigHelper.assertGdbClient(jobConfig);
ConfigHelper.assertLabels(jobConfig);
try {
exportType = Key.ExportType.valueOf(jobConfig.getString(Key.EXPORT_TYPE));
} catch (NullPointerException | IllegalArgumentException e) {
throw DataXException.asDataXException(GdbReaderErrorCode.BAD_CONFIG_VALUE, Key.EXPORT_TYPE);
}
}
@Override
public void prepare() {
/**
* 注意:此方法仅执行一次。
* 最佳实践:如果 Job 中有需要进行数据同步之前的处理,可以在此处完成,如果没有必要则可以直接去掉。
*/
try {
graph = new ScriptGdbGraph(jobConfig, exportType);
} catch (RuntimeException e) {
throw DataXException.asDataXException(GdbReaderErrorCode.FAIL_CLIENT_CONNECT, e.getMessage());
}
}
@Override
public List<Configuration> split(int adviceNumber) {
/**
* 注意:此方法仅执行一次。
* 最佳实践:通常采用工具静态类完成把 Job 配置切分成多个 Task 配置的工作。
* 这里的 adviceNumber 是框架根据用户的同步速度的要求建议的切分份数,仅供参考,不是强制必须切分的份数。
*/
List<String> labels = ConfigHelper.assertLabels(jobConfig);
/**
* 配置label列表为空时,尝试查询GDB中所有label,添加到读取列表
*/
if (labels.isEmpty()) {
try {
labels.addAll(graph.getLabels().keySet());
} catch (RuntimeException ex) {
throw DataXException.asDataXException(GdbReaderErrorCode.FAIL_FETCH_LABELS, ex.getMessage());
}
}
if (labels.isEmpty()) {
throw DataXException.asDataXException(GdbReaderErrorCode.FAIL_FETCH_LABELS, "none labels to read");
}
return ConfigHelper.splitConfig(jobConfig, labels);
}
@Override
public void post() {
/**
* 注意:此方法仅执行一次。
* 最佳实践:如果 Job 中有需要进行数据同步之后的后续处理,可以在此处完成。
*/
}
@Override
public void destroy() {
/**
* 注意:此方法仅执行一次。
* 最佳实践:通常配合 Job 中的 post() 方法一起完成 Job 的资源释放。
*/
try {
graph.close();
} catch (Exception ex) {
LOG.error("Failed to close client : {}", ex);
}
}
}
public static class Task extends Reader.Task {
private static final Logger LOG = LoggerFactory.getLogger(Task.class);
private static MappingRule rule;
private Configuration taskConfig;
private String fetchLabel = null;
private int rangeSplitSize;
private int fetchBatchSize;
@Override
public void init() {
this.taskConfig = super.getPluginJobConf();
/**
* 注意:此方法每个 Task 都会执行一次。
* 最佳实践:此处通过对 taskConfig 配置的读取,进而初始化一些资源为 startRead()做准备。
*/
fetchLabel = taskConfig.getString(Key.LABEL);
fetchBatchSize = taskConfig.getInt(Key.FETCH_BATCH_SIZE, DEFAULT_FETCH_BATCH_SIZE);
rangeSplitSize = taskConfig.getInt(Key.RANGE_SPLIT_SIZE, fetchBatchSize * 10);
rule = MappingRuleFactory.getInstance().create(taskConfig, exportType);
}
@Override
public void prepare() {
/**
* 注意:此方法仅执行一次。
* 最佳实践:如果 Job 中有需要进行数据同步之后的处理,可以在此处完成,如果没有必要则可以直接去掉。
*/
}
@Override
public void startRead(RecordSender recordSender) {
/**
* 注意:此方法每个 Task 都会执行一次。
* 最佳实践:此处适当封装确保简洁清晰完成数据读取工作。
*/
String start = "";
while (true) {
List<String> ids;
try {
ids = graph.fetchIds(fetchLabel, start, rangeSplitSize);
if (ids.isEmpty()) {
break;
}
start = ids.get(ids.size() - 1);
} catch (Exception ex) {
throw DataXException.asDataXException(GdbReaderErrorCode.FAIL_FETCH_IDS, ex.getMessage());
}
// send range fetch async
int count = ids.size();
List<ResultSet> resultSets = new LinkedList<>();
for (int pos = 0; pos < count; pos += fetchBatchSize) {
int rangeSize = Math.min(fetchBatchSize, count - pos);
String endId = ids.get(pos + rangeSize - 1);
String beginId = ids.get(pos);
List<String> propNames = rule.isHasProperty() ? rule.getPropertyNames() : null;
try {
resultSets.add(graph.fetchElementsAsync(fetchLabel, beginId, endId, propNames));
} catch (Exception ex) {
// just print error logs and continues
LOG.error("failed to request label: {}, start: {}, end: {}, e: {}", fetchLabel, beginId, endId, ex);
}
}
// get range fetch dsl results
resultSets.forEach(results -> {
try {
List<GdbElement> elements = graph.getElement(results);
elements.forEach(element -> {
Record record = recordSender.createRecord();
DefaultGdbMapper.getMapper(rule).accept(element, record);
recordSender.sendToWriter(record);
});
recordSender.flush();
} catch (Exception ex) {
LOG.error("failed to send records e {}", ex);
}
});
}
}
@Override
public void post() {
/**
* 注意:此方法每个 Task 都会执行一次。
* 最佳实践:如果 Task 中有需要进行数据同步之后的后续处理,可以在此处完成。
*/
}
@Override
public void destroy() {
/**
* 注意:此方法每个 Task 都会执行一次。
* 最佳实践:通常配合Task 中的 post() 方法一起完成 Task 的资源释放。
*/
}
}
}
@@ -0,0 +1,39 @@
package com.alibaba.datax.plugin.reader.gdbreader;
import com.alibaba.datax.common.spi.ErrorCode;
public enum GdbReaderErrorCode implements ErrorCode {
/**
*
*/
BAD_CONFIG_VALUE("GdbReader-00", "The value you configured is invalid."),
FAIL_CLIENT_CONNECT("GdbReader-02", "GDB connection is abnormal."),
UNSUPPORTED_TYPE("GdbReader-03", "Unsupported data type conversion."),
FAIL_FETCH_LABELS("GdbReader-04", "Error pulling all labels, it is recommended to configure the specified label pull."),
FAIL_FETCH_IDS("GdbReader-05", "Pull range id error."),
;
private final String code;
private final String description;
private GdbReaderErrorCode(String code, String description) {
this.code = code;
this.description = description;
}
@Override
public String getCode() {
return this.code;
}
@Override
public String getDescription() {
return this.description;
}
@Override
public String toString() {
return String.format("Code:[%s], Description:[%s]. ", this.code,
this.description);
}
}
@@ -0,0 +1,86 @@
package com.alibaba.datax.plugin.reader.gdbreader;
public final class Key {
/**
* 此处声明插件用到的需要插件使用者提供的配置项
*/
public final static String HOST = "host";
public final static String PORT = "port";
public final static String USERNAME = "username";
public static final String PASSWORD = "password";
public static final String LABEL = "labels";
public static final String EXPORT_TYPE = "labelType";
public static final String RANGE_SPLIT_SIZE = "RangeSplitSize";
public static final String FETCH_BATCH_SIZE = "fetchBatchSize";
public static final String COLUMN = "column";
public static final String COLUMN_NAME = "name";
public static final String COLUMN_TYPE = "type";
public static final String COLUMN_NODE_TYPE = "columnType";
public enum ExportType {
/**
* Import vertices
*/
VERTEX,
/**
* Import edges
*/
EDGE
}
public enum ColumnType {
/**
* vertex or edge id
*/
primaryKey,
/**
* vertex or edge label
*/
primaryLabel,
/**
* vertex property
*/
vertexProperty,
/**
* collects all vertex property to Json list
*/
vertexJsonProperty,
/**
* start vertex id of edge
*/
srcPrimaryKey,
/**
* start vertex label of edge
*/
srcPrimaryLabel,
/**
* end vertex id of edge
*/
dstPrimaryKey,
/**
* end vertex label of edge
*/
dstPrimaryLabel,
/**
* edge property
*/
edgeProperty,
/**
* collects all edge property to Json list
*/
edgeJsonProperty,
}
}
@@ -0,0 +1,150 @@
/*
* (C) 2019-present Alibaba Group Holding Limited.
*
* This program is free software; you can redistribute it and/or modify
* it under the terms of the GNU General Public License version 2 as
* published by the Free Software Foundation.
*/
package com.alibaba.datax.plugin.reader.gdbreader.mapping;
import com.alibaba.datax.common.element.Record;
import com.alibaba.datax.plugin.reader.gdbreader.model.GdbElement;
import org.apache.tinkerpop.gremlin.structure.util.reference.ReferenceProperty;
import org.apache.tinkerpop.gremlin.structure.util.reference.ReferenceVertexProperty;
import java.util.List;
import java.util.Map;
import java.util.function.BiConsumer;
import java.util.function.Function;
import java.util.stream.Collectors;
/**
* @author : Liu Jianping
* @date : 2019/9/6
*/
public class DefaultGdbMapper {
public static BiConsumer<GdbElement, Record> getMapper(MappingRule rule) {
return (gdbElement, record) -> rule.getColumns().forEach(columnMappingRule -> {
Object value = null;
ValueType type = columnMappingRule.getValueType();
String name = columnMappingRule.getName();
Map<String, Object> props = gdbElement.getProperties();
switch (columnMappingRule.getColumnType()) {
case dstPrimaryKey:
value = gdbElement.getTo();
break;
case srcPrimaryKey:
value = gdbElement.getFrom();
break;
case primaryKey:
value = gdbElement.getId();
break;
case primaryLabel:
value = gdbElement.getLabel();
break;
case dstPrimaryLabel:
value = gdbElement.getToLabel();
break;
case srcPrimaryLabel:
value = gdbElement.getFromLabel();
break;
case vertexProperty:
value = forVertexOnePropertyValue().apply(props.get(name));
break;
case edgeProperty:
value = forEdgePropertyValue().apply(props.get(name));
break;
case edgeJsonProperty:
value = forEdgeJsonProperties().apply(props);
break;
case vertexJsonProperty:
value = forVertexJsonProperties().apply(props);
break;
default:
break;
}
record.addColumn(type.applyObject(value));
});
}
/**
* parser ReferenceProperty value for edge
*
* @return property value
*/
private static Function<Object, Object> forEdgePropertyValue() {
return prop -> {
if (prop instanceof ReferenceProperty) {
return ((ReferenceProperty) prop).value();
}
return null;
};
}
/**
* parser ReferenceVertexProperty value for vertex
*
* @return the first property value in list
*/
private static Function<Object, Object> forVertexOnePropertyValue() {
return props -> {
if (props instanceof List<?>) {
// get the first one property if more than one
Object o = ((List) props).get(0);
if (o instanceof ReferenceVertexProperty) {
return ((ReferenceVertexProperty) o).value();
}
}
return null;
};
}
/**
* parser all edge properties to json string
*
* @return json string
*/
private static Function<Map<String, Object>, String> forEdgeJsonProperties() {
return props -> "{\"properties\":[" +
props.entrySet().stream().filter(p -> p.getValue() instanceof ReferenceProperty)
.map(p -> "{\"k\":\"" + ((ReferenceProperty) p.getValue()).key() + "\"," +
"\"t\":\"" + ((ReferenceProperty) p.getValue()).value().getClass().getSimpleName().toLowerCase() + "\"," +
"\"v\":\"" + String.valueOf(((ReferenceProperty) p.getValue()).value()) + "\"}")
.collect(Collectors.joining(",")) +
"]}";
}
/**
* parser all vertex properties to json string, include set-property
*
* @return json string
*/
private static Function<Map<String, Object>, String> forVertexJsonProperties() {
return props -> "{\"properties\":[" +
props.entrySet().stream().filter(p -> p.getValue() instanceof List<?>)
.map(p -> forVertexPropertyStr().apply((List<?>) p.getValue()))
.collect(Collectors.joining(",")) +
"]}";
}
/**
* parser one vertex property to json string item, set 'cardinality'
*
* @return json string item
*/
private static Function<List<?>, String> forVertexPropertyStr() {
return vp -> {
final String setFlag = vp.size() > 1 ? "set" : "single";
return vp.stream().filter(p -> p instanceof ReferenceVertexProperty)
.map(p -> "{\"k\":\"" + ((ReferenceVertexProperty) p).key() + "\"," +
"\"t\":\"" + ((ReferenceVertexProperty) p).value().getClass().getSimpleName().toLowerCase() + "\"," +
"\"v\":\"" + String.valueOf(((ReferenceVertexProperty) p).value()) + "\"," +
"\"c\":\"" + setFlag + "\"}")
.collect(Collectors.joining(","));
};
}
}
@@ -0,0 +1,79 @@
/*
* (C) 2019-present Alibaba Group Holding Limited.
*
* This program is free software; you can redistribute it and/or modify
* it under the terms of the GNU General Public License version 2 as
* published by the Free Software Foundation.
*/
package com.alibaba.datax.plugin.reader.gdbreader.mapping;
import com.alibaba.datax.common.exception.DataXException;
import com.alibaba.datax.plugin.reader.gdbreader.GdbReaderErrorCode;
import com.alibaba.datax.plugin.reader.gdbreader.Key.ColumnType;
import com.alibaba.datax.plugin.reader.gdbreader.Key.ExportType;
import lombok.Data;
import java.util.ArrayList;
import java.util.List;
/**
* @author : Liu Jianping
* @date : 2019/9/6
*/
@Data
public class MappingRule {
private boolean hasRelation = false;
private boolean hasProperty = false;
private ExportType type = ExportType.VERTEX;
/**
* property names for property key-value
*/
private List<String> propertyNames = new ArrayList<>();
private List<ColumnMappingRule> columns = new ArrayList<>();
void addColumn(ColumnType columnType, ValueType type, String name) {
ColumnMappingRule rule = new ColumnMappingRule();
rule.setColumnType(columnType);
rule.setName(name);
rule.setValueType(type);
if (columnType == ColumnType.vertexProperty || columnType == ColumnType.edgeProperty) {
propertyNames.add(name);
hasProperty = true;
}
boolean hasTo = columnType == ColumnType.dstPrimaryKey || columnType == ColumnType.dstPrimaryLabel;
boolean hasFrom = columnType == ColumnType.srcPrimaryKey || columnType == ColumnType.srcPrimaryLabel;
if (hasTo || hasFrom) {
hasRelation = true;
}
columns.add(rule);
}
void addJsonColumn(ColumnType columnType) {
ColumnMappingRule rule = new ColumnMappingRule();
rule.setColumnType(columnType);
rule.setName("json");
rule.setValueType(ValueType.STRING);
if (!propertyNames.isEmpty()) {
throw DataXException.asDataXException(GdbReaderErrorCode.BAD_CONFIG_VALUE, "JsonProperties should be only property");
}
columns.add(rule);
hasProperty = true;
}
@Data
protected static class ColumnMappingRule {
private String name = null;
private ValueType valueType = null;
private ColumnType columnType = null;
}
}
@@ -0,0 +1,76 @@
/*
* (C) 2019-present Alibaba Group Holding Limited.
*
* This program is free software; you can redistribute it and/or modify
* it under the terms of the GNU General Public License version 2 as
* published by the Free Software Foundation.
*/
package com.alibaba.datax.plugin.reader.gdbreader.mapping;
import com.alibaba.datax.common.exception.DataXException;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.reader.gdbreader.GdbReaderErrorCode;
import com.alibaba.datax.plugin.reader.gdbreader.Key;
import com.alibaba.datax.plugin.reader.gdbreader.Key.ColumnType;
import com.alibaba.datax.plugin.reader.gdbreader.Key.ExportType;
import com.alibaba.datax.plugin.reader.gdbreader.util.ConfigHelper;
import java.util.List;
/**
* @author : Liu Jianping
* @date : 2019/9/20
*/
public class MappingRuleFactory {
private static final MappingRuleFactory instance = new MappingRuleFactory();
public static MappingRuleFactory getInstance() {
return instance;
}
public MappingRule create(Configuration config, ExportType exportType) {
MappingRule rule = new MappingRule();
rule.setType(exportType);
List<Configuration> configurationList = config.getListConfiguration(Key.COLUMN);
for (Configuration column : configurationList) {
ColumnType columnType;
try {
columnType = ColumnType.valueOf(column.getString(Key.COLUMN_NODE_TYPE));
} catch (NullPointerException | IllegalArgumentException e) {
throw DataXException.asDataXException(GdbReaderErrorCode.BAD_CONFIG_VALUE, Key.COLUMN_NODE_TYPE);
}
if (exportType == ExportType.VERTEX) {
// only id/label/property column allow when vertex
ConfigHelper.assertConfig(Key.COLUMN_NODE_TYPE, () ->
columnType == ColumnType.primaryKey || columnType == ColumnType.primaryLabel
|| columnType == ColumnType.vertexProperty || columnType == ColumnType.vertexJsonProperty);
} else if (exportType == ExportType.EDGE) {
// edge
ConfigHelper.assertConfig(Key.COLUMN_NODE_TYPE, () ->
columnType == ColumnType.primaryKey || columnType == ColumnType.primaryLabel
|| columnType == ColumnType.srcPrimaryKey || columnType == ColumnType.srcPrimaryLabel
|| columnType == ColumnType.dstPrimaryKey || columnType == ColumnType.dstPrimaryLabel
|| columnType == ColumnType.edgeProperty || columnType == ColumnType.edgeJsonProperty);
}
if (columnType == ColumnType.edgeProperty || columnType == ColumnType.vertexProperty) {
String name = column.getString(Key.COLUMN_NAME);
ValueType propType = ValueType.fromShortName(column.getString(Key.COLUMN_TYPE));
ConfigHelper.assertConfig(Key.COLUMN_NAME, () -> name != null);
if (propType == null) {
throw DataXException.asDataXException(GdbReaderErrorCode.UNSUPPORTED_TYPE, Key.COLUMN_TYPE);
}
rule.addColumn(columnType, propType, name);
} else if (columnType == ColumnType.vertexJsonProperty || columnType == ColumnType.edgeJsonProperty) {
rule.addJsonColumn(columnType);
} else {
rule.addColumn(columnType, ValueType.STRING, null);
}
}
return rule;
}
}
@@ -0,0 +1,128 @@
/*
* (C) 2019-present Alibaba Group Holding Limited.
*
* This program is free software; you can redistribute it and/or modify
* it under the terms of the GNU General Public License version 2 as
* published by the Free Software Foundation.
*/
package com.alibaba.datax.plugin.reader.gdbreader.mapping;
import com.alibaba.datax.common.element.BoolColumn;
import com.alibaba.datax.common.element.Column;
import com.alibaba.datax.common.element.DoubleColumn;
import com.alibaba.datax.common.element.LongColumn;
import com.alibaba.datax.common.element.StringColumn;
import java.util.HashMap;
import java.util.Map;
import java.util.function.Function;
/**
* @author : Liu Jianping
* @date : 2019/9/6
*/
public enum ValueType {
/**
* transfer gdb element object value to DataX Column data
* <p>
* int, long -> LongColumn
* float, double -> DoubleColumn
* bool -> BooleanColumn
* string -> StringColumn
*/
INT(Integer.class, "int", ValueTypeHolder::longColumnMapper),
INTEGER(Integer.class, "integer", ValueTypeHolder::longColumnMapper),
LONG(Long.class, "long", ValueTypeHolder::longColumnMapper),
DOUBLE(Double.class, "double", ValueTypeHolder::doubleColumnMapper),
FLOAT(Float.class, "float", ValueTypeHolder::doubleColumnMapper),
BOOLEAN(Boolean.class, "boolean", ValueTypeHolder::boolColumnMapper),
STRING(String.class, "string", ValueTypeHolder::stringColumnMapper),
;
private Class<?> type = null;
private String shortName = null;
private Function<Object, Column> columnFunc = null;
ValueType(Class<?> type, String name, Function<Object, Column> columnFunc) {
this.type = type;
this.shortName = name;
this.columnFunc = columnFunc;
ValueTypeHolder.shortName2type.put(shortName, this);
}
public static ValueType fromShortName(String name) {
return ValueTypeHolder.shortName2type.get(name);
}
public Column applyObject(Object value) {
if (value == null) {
return null;
}
return columnFunc.apply(value);
}
private static class ValueTypeHolder {
private static Map<String, ValueType> shortName2type = new HashMap<>();
private static LongColumn longColumnMapper(Object o) {
long v;
if (o instanceof Integer) {
v = (int) o;
} else if (o instanceof Long) {
v = (long) o;
} else if (o instanceof String) {
v = Long.valueOf((String) o);
} else {
throw new RuntimeException("Failed to cast " + o.getClass() + " to Long");
}
return new LongColumn(v);
}
private static DoubleColumn doubleColumnMapper(Object o) {
double v;
if (o instanceof Integer) {
v = (double) (int) o;
} else if (o instanceof Long) {
v = (double) (long) o;
} else if (o instanceof Float) {
v = (double) (float) o;
} else if (o instanceof Double) {
v = (double) o;
} else if (o instanceof String) {
v = Double.valueOf((String) o);
} else {
throw new RuntimeException("Failed to cast " + o.getClass() + " to Double");
}
return new DoubleColumn(v);
}
private static BoolColumn boolColumnMapper(Object o) {
boolean v;
if (o instanceof Integer) {
v = ((int) o != 0);
} else if (o instanceof Long) {
v = ((long) o != 0);
} else if (o instanceof Boolean) {
v = (boolean) o;
} else if (o instanceof String) {
v = Boolean.valueOf((String) o);
} else {
throw new RuntimeException("Failed to cast " + o.getClass() + " to Boolean");
}
return new BoolColumn(v);
}
private static StringColumn stringColumnMapper(Object o) {
if (o instanceof String) {
return new StringColumn((String) o);
} else {
return new StringColumn(String.valueOf(o));
}
}
}
}
@@ -0,0 +1,89 @@
/*
* (C) 2019-present Alibaba Group Holding Limited.
*
* This program is free software; you can redistribute it and/or modify
* it under the terms of the GNU General Public License version 2 as
* published by the Free Software Foundation.
*/
package com.alibaba.datax.plugin.reader.gdbreader.model;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.reader.gdbreader.Key;
import org.apache.tinkerpop.gremlin.driver.Client;
import org.apache.tinkerpop.gremlin.driver.Cluster;
import org.apache.tinkerpop.gremlin.driver.RequestOptions;
import org.apache.tinkerpop.gremlin.driver.Result;
import org.apache.tinkerpop.gremlin.driver.ResultSet;
import org.apache.tinkerpop.gremlin.driver.ser.Serializers;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;
/**
* @author : Liu Jianping
* @date : 2019/9/6
*/
public abstract class AbstractGdbGraph implements GdbGraph {
final static int DEFAULT_TIMEOUT = 30000;
private static final Logger log = LoggerFactory.getLogger(AbstractGdbGraph.class);
private Client client;
AbstractGdbGraph() {
}
AbstractGdbGraph(Configuration config) {
log.info("init graphdb client");
String host = config.getString(Key.HOST);
int port = config.getInt(Key.PORT);
String username = config.getString(Key.USERNAME);
String password = config.getString(Key.PASSWORD);
try {
Cluster cluster = Cluster.build(host).port(port).credentials(username, password)
.serializer(Serializers.GRAPHBINARY_V1D0)
.maxContentLength(1024 * 1024)
.resultIterationBatchSize(64)
.create();
client = cluster.connect().init();
warmClient();
} catch (RuntimeException e) {
log.error("Failed to connect to GDB {}:{}, due to {}", host, port, e);
throw e;
}
}
protected List<Result> runInternal(String dsl, Map<String, Object> params) throws Exception {
return runInternalAsync(dsl, params).all().get(DEFAULT_TIMEOUT + 1000, TimeUnit.MILLISECONDS);
}
protected ResultSet runInternalAsync(String dsl, Map<String, Object> params) throws Exception {
RequestOptions.Builder options = RequestOptions.build().timeout(DEFAULT_TIMEOUT);
if (params != null && !params.isEmpty()) {
params.forEach(options::addParameter);
}
return client.submitAsync(dsl, options.create()).get(DEFAULT_TIMEOUT, TimeUnit.MILLISECONDS);
}
private void warmClient() {
try {
runInternal("g.V('test')", null);
log.info("warm graphdb client over");
} catch (Exception e) {
log.error("warmClient error");
throw new RuntimeException(e);
}
}
@Override
public void close() throws Exception {
if (client != null) {
log.info("close graphdb client");
client.close();
}
}
}
@@ -0,0 +1,39 @@
/*
* (C) 2019-present Alibaba Group Holding Limited.
*
* This program is free software; you can redistribute it and/or modify
* it under the terms of the GNU General Public License version 2 as
* published by the Free Software Foundation.
*/
package com.alibaba.datax.plugin.reader.gdbreader.model;
import lombok.Data;
import java.util.HashMap;
import java.util.Map;
/**
* @author : Liu Jianping
* @date : 2019/9/6
*/
@Data
public class GdbElement {
String id = null;
String label = null;
String to = null;
String from = null;
String toLabel = null;
String fromLabel = null;
Map<String, Object> properties = new HashMap<>();
public GdbElement() {
}
public GdbElement(String id, String label) {
this.id = id;
this.label = label;
}
}
@@ -0,0 +1,65 @@
/*
* (C) 2019-present Alibaba Group Holding Limited.
*
* This program is free software; you can redistribute it and/or modify
* it under the terms of the GNU General Public License version 2 as
* published by the Free Software Foundation.
*/
package com.alibaba.datax.plugin.reader.gdbreader.model;
import org.apache.tinkerpop.gremlin.driver.ResultSet;
import java.util.List;
import java.util.Map;
/**
* @author : Liu Jianping
* @date : 2019/9/6
*/
public interface GdbGraph extends AutoCloseable {
/**
* Get All labels of GraphDB
*
* @return labels map included numbers
*/
Map<String, Long> getLabels();
/**
* Get the Ids list of special 'label', size up to 'limit'
*
* @param label is Label of Vertex or Edge
* @param start of Ids range to get
* @param limit size of Ids list
* @return Ids list
*/
List<String> fetchIds(String label, String start, long limit);
/**
* Fetch element in async mode, just send query dsl to server
*
* @param label node label to filter
* @param start range begin(included)
* @param end range end(included)
* @param propNames propKey list to fetch
* @return future to get result later
*/
ResultSet fetchElementsAsync(String label, String start, String end, List<String> propNames);
/**
* Get get element from Response @{ResultSet}
*
* @param results Response of Server
* @return element sets
*/
List<GdbElement> getElement(ResultSet results);
/**
* close graph client
*
* @throws Exception if fails
*/
@Override
void close() throws Exception;
}
@@ -0,0 +1,192 @@
/*
* (C) 2019-present Alibaba Group Holding Limited.
*
* This program is free software; you can redistribute it and/or modify
* it under the terms of the GNU General Public License version 2 as
* published by the Free Software Foundation.
*/
package com.alibaba.datax.plugin.reader.gdbreader.model;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.reader.gdbreader.Key.ExportType;
import org.apache.tinkerpop.gremlin.driver.Result;
import org.apache.tinkerpop.gremlin.driver.ResultSet;
import org.apache.tinkerpop.gremlin.structure.util.reference.ReferenceEdge;
import org.apache.tinkerpop.gremlin.structure.util.reference.ReferenceVertex;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;
/**
* @author : Liu Jianping
* @date : 2019/9/6
*/
public class ScriptGdbGraph extends AbstractGdbGraph {
private static final Logger log = LoggerFactory.getLogger(ScriptGdbGraph.class);
private final static String LABEL = "GDB___LABEL";
private final static String START_ID = "GDB___ID";
private final static String END_ID = "GDB___ID_END";
private final static String LIMIT = "GDB___LIMIT";
private final static String FETCH_VERTEX_IDS_DSL = "g.V().hasLabel(" + LABEL + ").has(id, gt(" + START_ID + ")).limit(" + LIMIT + ").id()";
private final static String FETCH_EDGE_IDS_DSL = "g.E().hasLabel(" + LABEL + ").has(id, gt(" + START_ID + ")).limit(" + LIMIT + ").id()";
private final static String FETCH_VERTEX_LABELS_DSL = "g.V().groupCount().by(label)";
private final static String FETCH_EDGE_LABELS_DSL = "g.E().groupCount().by(label)";
/**
* fetch node range [START_ID, END_ID]
*/
private final static String FETCH_RANGE_VERTEX_DSL = "g.V().hasLabel(" + LABEL + ").has(id, gte(" + START_ID + ")).has(id, lte(" + END_ID + "))";
private final static String FETCH_RANGE_EDGE_DSL = "g.E().hasLabel(" + LABEL + ").has(id, gte(" + START_ID + ")).has(id, lte(" + END_ID + "))";
private final static String PART_WITH_PROP_DSL = ".as('a').project('node', 'props').by(select('a')).by(select('a').propertyMap(";
private final ExportType exportType;
public ScriptGdbGraph(ExportType exportType) {
super();
this.exportType = exportType;
}
public ScriptGdbGraph(Configuration config, ExportType exportType) {
super(config);
this.exportType = exportType;
}
@Override
public List<String> fetchIds(final String label, final String start, long limit) {
Map<String, Object> params = new HashMap<String, Object>(3) {{
put(LABEL, label);
put(START_ID, start);
put(LIMIT, limit);
}};
String fetchDsl = exportType == ExportType.VERTEX ? FETCH_VERTEX_IDS_DSL : FETCH_EDGE_IDS_DSL;
List<String> ids = new ArrayList<>();
try {
List<Result> results = runInternal(fetchDsl, params);
// transfer result to id string
results.forEach(id -> ids.add(id.getString()));
} catch (Exception e) {
log.error("fetch range node failed, label {}, start {}", label, start);
throw new RuntimeException(e);
}
return ids;
}
@Override
public ResultSet fetchElementsAsync(final String label, final String start, final String end, final List<String> propNames) {
Map<String, Object> params = new HashMap<>(3);
params.put(LABEL, label);
params.put(START_ID, start);
params.put(END_ID, end);
String prefixDsl = exportType == ExportType.VERTEX ? FETCH_RANGE_VERTEX_DSL : FETCH_RANGE_EDGE_DSL;
StringBuilder fetchDsl = new StringBuilder(prefixDsl);
if (propNames != null) {
fetchDsl.append(PART_WITH_PROP_DSL);
for (int i = 0; i < propNames.size(); i++) {
String propName = "GDB___PK" + String.valueOf(i);
params.put(propName, propNames.get(i));
fetchDsl.append(propName);
if (i != propNames.size() - 1) {
fetchDsl.append(", ");
}
}
fetchDsl.append("))");
}
try {
return runInternalAsync(fetchDsl.toString(), params);
} catch (Exception e) {
log.error("Failed to fetch range node startId {}, end {} , e {}", start, end, e);
throw new RuntimeException(e);
}
}
@Override
@SuppressWarnings("unchecked")
public List<GdbElement> getElement(ResultSet results) {
List<GdbElement> elements = new LinkedList<>();
try {
List<Result> resultList = results.all().get(DEFAULT_TIMEOUT + 1000, TimeUnit.MILLISECONDS);
resultList.forEach(n -> {
Object o = n.getObject();
GdbElement element = new GdbElement();
if (o instanceof Map) {
// project response
Object node = ((Map) o).get("node");
Object props = ((Map) o).get("props");
mapNodeToElement(node, element);
mapPropToElement((Map<String, Object>) props, element);
} else {
// range node response
mapNodeToElement(n.getObject(), element);
}
if (element.getId() != null) {
elements.add(element);
}
});
} catch (Exception e) {
log.error("Failed to get node: {}", e);
throw new RuntimeException(e);
}
return elements;
}
private void mapNodeToElement(Object node, GdbElement element) {
if (node instanceof ReferenceVertex) {
ReferenceVertex v = (ReferenceVertex) node;
element.setId((String) v.id());
element.setLabel(v.label());
} else if (node instanceof ReferenceEdge) {
ReferenceEdge e = (ReferenceEdge) node;
element.setId((String) e.id());
element.setLabel(e.label());
element.setTo((String) e.inVertex().id());
element.setToLabel(e.inVertex().label());
element.setFrom((String) e.outVertex().id());
element.setFromLabel(e.outVertex().label());
}
}
private void mapPropToElement(Map<String, Object> props, GdbElement element) {
element.setProperties(props);
}
@Override
public Map<String, Long> getLabels() {
String dsl = exportType == ExportType.VERTEX ? FETCH_VERTEX_LABELS_DSL : FETCH_EDGE_LABELS_DSL;
try {
List<Result> results = runInternal(dsl, null);
Map<String, Long> labelMap = new HashMap<>(2);
Map<?, ?> labels = results.get(0).get(Map.class);
labels.forEach((k, v) -> {
String label = (String) k;
Long count = (Long) v;
labelMap.put(label, count);
});
return labelMap;
} catch (Exception e) {
log.error("Failed to fetch label list, please give special labels and run again, e {}", e);
throw new RuntimeException(e);
}
}
}
@@ -0,0 +1,77 @@
/*
* (C) 2019-present Alibaba Group Holding Limited.
*
* This program is free software; you can redistribute it and/or modify
* it under the terms of the GNU General Public License version 2 as
* published by the Free Software Foundation.
*/
package com.alibaba.datax.plugin.reader.gdbreader.util;
import com.alibaba.datax.common.exception.DataXException;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.reader.gdbreader.GdbReaderErrorCode;
import com.alibaba.datax.plugin.reader.gdbreader.Key;
import org.apache.commons.lang3.StringUtils;
import java.io.IOException;
import java.io.InputStream;
import java.util.ArrayList;
import java.util.List;
import java.util.function.Supplier;
/**
* @author : Liu Jianping
* @date : 2019/9/6
*/
public interface ConfigHelper {
static void assertConfig(String key, Supplier<Boolean> f) {
if (!f.get()) {
throw DataXException.asDataXException(GdbReaderErrorCode.BAD_CONFIG_VALUE, key);
}
}
static void assertHasContent(Configuration config, String key) {
assertConfig(key, () -> StringUtils.isNotBlank(config.getString(key)));
}
static void assertGdbClient(Configuration config) {
assertHasContent(config, Key.HOST);
assertConfig(Key.PORT, () -> config.getInt(Key.PORT) > 0);
assertHasContent(config, Key.USERNAME);
assertHasContent(config, Key.PASSWORD);
}
static List<String> assertLabels(Configuration config) {
Object labels = config.get(Key.LABEL);
if (!(labels instanceof List)) {
throw DataXException.asDataXException(GdbReaderErrorCode.BAD_CONFIG_VALUE, "labels should be List");
}
List<?> list = (List<?>) labels;
List<String> configLabels = new ArrayList<>(0);
list.forEach(n -> configLabels.add(String.valueOf(n)));
return configLabels;
}
static List<Configuration> splitConfig(Configuration config, List<String> labels) {
List<Configuration> configs = new ArrayList<>();
for (String label : labels) {
Configuration conf = config.clone();
conf.set(Key.LABEL, label);
configs.add(conf);
}
return configs;
}
static Configuration fromClasspath(String name) {
try (InputStream is = Thread.currentThread().getContextClassLoader().getResourceAsStream(name)) {
return Configuration.from(is);
} catch (IOException e) {
throw new IllegalArgumentException("File not found: " + name);
}
}
}
+6
View File
@@ -0,0 +1,6 @@
{
"name": "gdbreader",
"class": "com.alibaba.datax.plugin.reader.gdbreader.GdbReader",
"description": "useScene: prod. mechanism: connect GDB with gremlin-client, execute 'g.V().propertyMap() or g.E().propertyMap()' to get record",
"developer": "alibaba"
}
@@ -0,0 +1,77 @@
{
"job": {
"setting": {
"speed": {
"channel": 1
},
"errorLimit": {
"record": 1
}
},
"content": [
{
"reader": {
"name": "gdbreader",
"parameter": {
"host": "10.218.145.24",
"port": 8182,
"username": "***",
"password": "***",
"labelType": "EDGE",
"labels": ["label1", "label2"],
"column": [
{
"name": "id",
"type": "string",
"columnType": "primaryKey"
},
{
"name": "label",
"type": "string",
"columnType": "primaryLabel"
},
{
"name": "srcId",
"type": "string",
"columnType": "srcPrimaryKey"
},
{
"name": "srcLabel",
"type": "string",
"columnType": "srcPrimaryLabel"
},
{
"name": "dstId",
"type": "string",
"columnType": "srcPrimaryKey"
},
{
"name": "dstLabel",
"type": "string",
"columnType": "srcPrimaryLabel"
},
{
"name": "name",
"type": "string",
"columnType": "edgeProperty"
},
{
"name": "weight",
"type": "double",
"columnType": "edgeProperty"
}
]
}
},
"writer": {
"name": "streamwriter",
"parameter": {
"print": true
}
}
}
]
}
}
+45 -15
View File
@@ -41,6 +41,14 @@ GDBWriter通过DataX框架获取Reader生成的协议数据,使用`g.addV/E(GD
{
"random": "60,64",
"type": "string"
},
{
"random": "100,1000",
"type": "long"
},
{
"random": "32,48",
"type": "string"
}
],
"sliceRecordCount": 1000
@@ -55,20 +63,32 @@ GDBWriter通过DataX框架获取Reader生成的协议数据,使用`g.addV/E(GD
"password": "***",
"writeMode": "INSERT",
"labelType": "VERTEX",
"label": "${1}",
"label": "#{1}",
"idTransRule": "none",
"session": true,
"maxRecordsInBatch": 64,
"column": [
{
"name": "id",
"value": "${0}",
"value": "#{0}",
"type": "string",
"columnType": "primaryKey"
},
{
"name": "vertex_propKey",
"value": "${2}",
"value": "#{2}",
"type": "string",
"columnType": "vertexSetProperty"
},
{
"name": "vertex_propKey",
"value": "#{3}",
"type": "long",
"columnType": "vertexSetProperty"
},
{
"name": "vertex_propKey2",
"value": "#{4}",
"type": "string",
"columnType": "vertexProperty"
}
@@ -134,7 +154,7 @@ GDBWriter通过DataX框架获取Reader生成的协议数据,使用`g.addV/E(GD
"password": "***",
"writeMode": "INSERT",
"labelType": "EDGE",
"label": "${3}",
"label": "#{3}",
"idTransRule": "none",
"srcIdTransRule": "labelPrefix",
"dstIdTransRule": "labelPrefix",
@@ -144,25 +164,25 @@ GDBWriter通过DataX框架获取Reader生成的协议数据,使用`g.addV/E(GD
"column": [
{
"name": "id",
"value": "${0}",
"value": "#{0}",
"type": "string",
"columnType": "primaryKey"
},
{
"name": "id",
"value": "${1}",
"value": "#{1}",
"type": "string",
"columnType": "srcPrimaryKey"
},
{
"name": "id",
"value": "${2}",
"value": "#{2}",
"type": "string",
"columnType": "dstPrimaryKey"
},
{
"name": "edge_propKey",
"value": "${4}",
"value": "#{4}",
"type": "string",
"columnType": "edgeProperty"
}
@@ -199,7 +219,7 @@ GDBWriter通过DataX框架获取Reader生成的协议数据,使用`g.addV/E(GD
* 默认值:无
* **label**
* 描述:类型名,即点/边名称; label支持从源列中读取,如${0},表示取第一列字段作为label名。源列索引从0开始;
* 描述:类型名,即点/边名称; label支持从源列中读取,如#{0},表示取第一列字段作为label名。源列索引从0开始;
* 必选:是
* 默认值:无
@@ -211,12 +231,12 @@ GDBWriter通过DataX框架获取Reader生成的协议数据,使用`g.addV/E(GD
* 默认值:无
* **srcLabel**
* 描述:当label为边时,表示起点的点名称;srcLabel支持从源列中读取,如${0},表示取第一列字段作为label名。源列索引从0开始;
* 描述:当label为边时,表示起点的点名称;srcLabel支持从源列中读取,如#{0},表示取第一列字段作为label名。源列索引从0开始;
* 必选:labelType为边,srcIdTransRule为none时可不填写,否则必填;
* 默认值:无
* **dstLabel**
* 描述:当label为边时,表示终点的点名称;dstLabel支持从源列中读取,如${0},表示取第一列字段作为label名。源列索引从0开始;
* 描述:当label为边时,表示终点的点名称;dstLabel支持从源列中读取,如#{0},表示取第一列字段作为label名。源列索引从0开始;
* 必选:labelType为边,dstIdTransRule为none时可不填写,否则必填;
* 默认值:无
@@ -271,9 +291,9 @@ GDBWriter通过DataX框架获取Reader生成的协议数据,使用`g.addV/E(GD
* **column -> value**
* 描述:点/边映射关系的字段值;
* ${N}表示直接映射源端值,N为源端column索引,从0开始;${0}表示映射源端column第1个字段;
* test-${0} 表示源端值做拼接转换,${0}值前/后可添加固定字符串;
* ${0}-${1}表示做多字段拼接,也可在任意位置添加固定字符串,如test-${0}-test1-${1}-test2
* #{N}表示直接映射源端值,N为源端column索引,从0开始;#{0}表示映射源端column第1个字段;
* test-#{0} 表示源端值做拼接转换,#{0}值前/后可添加固定字符串;
* #{0}-#{1}表示做多字段拼接,也可在任意位置添加固定字符串,如test-#{0}-test1-#{1}-test2
* 必选:是
* 默认值:无
@@ -290,6 +310,7 @@ GDBWriter通过DataX框架获取Reader生成的协议数据,使用`g.addV/E(GD
* primaryKey:表示该字段是主键id
* 点枚举值:
* vertexPropertylabelType为点时,表示该字段是点的普通属性
* vertexSetPropertylabelType为点时,表示该字段是点的SET属性,value是SET属性中的一个属性值
* vertexJsonPropertylabelType为点时,表示是点json属性,value结构请见备注**json properties示例**,点配置最多只允许出现一个json属性;
* 边枚举值:
* srcPrimaryKeylabelType为边时,表示该字段是起点主键id
@@ -305,6 +326,14 @@ GDBWriter通过DataX框架获取Reader生成的协议数据,使用`g.addV/E(GD
> {"k":"age","t":"int","v":"20"},
> {"k":"sex","t":"string","v":"male"}
> ]}
>
> # json格式同样支持给点添加SET属性,格式如下
> {"properties":[
> {"k":"name","t":"string","v":"tom","c":"set"},
> {"k":"name","t":"string","v":"jack","c":"set"},
> {"k":"age","t":"int","v":"20"},
> {"k":"sex","t":"string","v":"male"}
> ]}
> ```
## 4 性能报告
@@ -367,4 +396,5 @@ DataX压测机器
- GDBWriter插件与用户查询DSL使用相同的GDB实例端口,导入时可能会影响查询性能
## FAQ
1. 使用SET属性需要升级GDB实例到`1.0.20`版本及以上。
2. 边只支持普通单值属性,不能给边写SET属性数据。
@@ -1,10 +1,5 @@
package com.alibaba.datax.plugin.writer.gdbwriter;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.*;
import java.util.function.Function;
import com.alibaba.datax.common.element.Record;
import com.alibaba.datax.common.exception.DataXException;
import com.alibaba.datax.common.plugin.RecordReceiver;
@@ -18,24 +13,33 @@ import com.alibaba.datax.plugin.writer.gdbwriter.mapping.MappingRule;
import com.alibaba.datax.plugin.writer.gdbwriter.mapping.MappingRuleFactory;
import com.alibaba.datax.plugin.writer.gdbwriter.model.GdbElement;
import com.alibaba.datax.plugin.writer.gdbwriter.model.GdbGraph;
import groovy.lang.Tuple2;
import io.netty.util.concurrent.DefaultThreadFactory;
import lombok.extern.slf4j.Slf4j;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class GdbWriter extends Writer {
private static final Logger log = LoggerFactory.getLogger(GdbWriter.class);
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
import java.util.concurrent.LinkedBlockingDeque;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.function.Function;
private static Function<Record, GdbElement> mapper = null;
private static GdbGraph globalGraph = null;
private static boolean session = false;
public class GdbWriter extends Writer {
private static final Logger log = LoggerFactory.getLogger(GdbWriter.class);
private static Function<Record, GdbElement> mapper = null;
private static GdbGraph globalGraph = null;
private static boolean session = false;
/**
* Job 中的方法仅执行一次,Task 中方法会由框架启动多个 Task 线程并行执行。
* <p/>
* 整个 Writer 执行流程是:
*
* <pre>
* Job类init-->prepare-->split
*
@@ -46,17 +50,16 @@ public class GdbWriter extends Writer {
* </pre>
*/
public static class Job extends Writer.Job {
private static final Logger LOG = LoggerFactory
.getLogger(Job.class);
private static final Logger LOG = LoggerFactory.getLogger(Job.class);
private Configuration jobConfig = null;
@Override
public void init() {
LOG.info("GDB datax plugin writer job init begin ...");
this.jobConfig = getPluginJobConf();
GdbWriterConfig.of(this.jobConfig);
LOG.info("GDB datax plugin writer job init end.");
LOG.info("GDB datax plugin writer job init begin ...");
this.jobConfig = getPluginJobConf();
GdbWriterConfig.of(this.jobConfig);
LOG.info("GDB datax plugin writer job init end.");
/**
* 注意:此方法仅执行一次。
@@ -71,37 +74,37 @@ public class GdbWriter extends Writer {
* 注意:此方法仅执行一次。
* 最佳实践:如果 Job 中有需要进行数据同步之前的处理,可以在此处完成,如果没有必要则可以直接去掉。
*/
super.prepare();
super.prepare();
MappingRule rule = MappingRuleFactory.getInstance().createV2(jobConfig);
final MappingRule rule = MappingRuleFactory.getInstance().createV2(this.jobConfig);
mapper = new DefaultGdbMapper().getMapper(rule);
session = jobConfig.getBool(Key.SESSION_STATE, false);
mapper = new DefaultGdbMapper(this.jobConfig).getMapper(rule);
session = this.jobConfig.getBool(Key.SESSION_STATE, false);
/**
* client connect check before task
*/
try {
globalGraph = GdbGraphManager.instance().getGraph(jobConfig, false);
} catch (RuntimeException e) {
globalGraph = GdbGraphManager.instance().getGraph(this.jobConfig, false);
} catch (final RuntimeException e) {
throw DataXException.asDataXException(GdbWriterErrorCode.FAIL_CLIENT_CONNECT, e.getMessage());
}
}
@Override
public List<Configuration> split(int mandatoryNumber) {
public List<Configuration> split(final int mandatoryNumber) {
/**
* 注意:此方法仅执行一次。
* 最佳实践:通常采用工具静态类完成把 Job 配置切分成多个 Task 配置的工作。
* 这里的 mandatoryNumber 是强制必须切分的份数。
*/
LOG.info("split begin...");
List<Configuration> configurationList = new ArrayList<Configuration>();
for (int i = 0; i < mandatoryNumber; i++) {
configurationList.add(this.jobConfig.clone());
}
LOG.info("split end...");
return configurationList;
LOG.info("split begin...");
final List<Configuration> configurationList = new ArrayList<Configuration>();
for (int i = 0; i < mandatoryNumber; i++) {
configurationList.add(this.jobConfig.clone());
}
LOG.info("split end...");
return configurationList;
}
@Override
@@ -127,7 +130,7 @@ public class GdbWriter extends Writer {
public static class Task extends Writer.Task {
private Configuration taskConfig;
private int failed = 0;
private int batchRecords;
private ExecutorService submitService = null;
@@ -139,24 +142,24 @@ public class GdbWriter extends Writer {
* 注意:此方法每个 Task 都会执行一次。
* 最佳实践:此处通过对 taskConfig 配置的读取,进而初始化一些资源为 startWrite()做准备。
*/
this.taskConfig = super.getPluginJobConf();
batchRecords = taskConfig.getInt(Key.MAX_RECORDS_IN_BATCH, GdbWriterConfig.DEFAULT_RECORD_NUM_IN_BATCH);
submitService = new ThreadPoolExecutor(1, 1, 0L,
TimeUnit.MILLISECONDS, new LinkedBlockingDeque<>(), new DefaultThreadFactory("submit-dsl"));
this.taskConfig = super.getPluginJobConf();
this.batchRecords = this.taskConfig.getInt(Key.MAX_RECORDS_IN_BATCH, GdbWriterConfig.DEFAULT_RECORD_NUM_IN_BATCH);
this.submitService = new ThreadPoolExecutor(1, 1, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingDeque<>(),
new DefaultThreadFactory("submit-dsl"));
if (!session) {
graph = globalGraph;
} else {
/**
* 分批创建session client,由于服务端groovy编译性能的限制
*/
try {
Thread.sleep((getTaskId()/10)*10000);
} catch (Exception e) {
// ...
}
graph = GdbGraphManager.instance().getGraph(taskConfig, session);
}
if (!session) {
this.graph = globalGraph;
} else {
/**
* 分批创建session client,由于服务端groovy编译性能的限制
*/
try {
Thread.sleep((getTaskId() / 10) * 10000);
} catch (final Exception e) {
// ...
}
this.graph = GdbGraphManager.instance().getGraph(this.taskConfig, session);
}
}
@Override
@@ -165,64 +168,69 @@ public class GdbWriter extends Writer {
* 注意:此方法每个 Task 都会执行一次。
* 最佳实践:如果 Task 中有需要进行数据同步之前的处理,可以在此处完成,如果没有必要则可以直接去掉。
*/
super.prepare();
super.prepare();
}
@Override
public void startWrite(RecordReceiver recordReceiver) {
public void startWrite(final RecordReceiver recordReceiver) {
/**
* 注意:此方法每个 Task 都会执行一次。
* 最佳实践:此处适当封装确保简洁清晰完成数据写入工作。
*/
Record r;
Future<Boolean> future = null;
List<Tuple2<Record, GdbElement>> records = new ArrayList<>(batchRecords);
Record r;
Future<Boolean> future = null;
List<Tuple2<Record, GdbElement>> records = new ArrayList<>(this.batchRecords);
while ((r = recordReceiver.getFromReader()) != null) {
records.add(new Tuple2<>(r, mapper.apply(r)));
while ((r = recordReceiver.getFromReader()) != null) {
try {
records.add(new Tuple2<>(r, mapper.apply(r)));
} catch (final Exception ex) {
getTaskPluginCollector().collectDirtyRecord(r, ex);
continue;
}
if (records.size() >= batchRecords) {
wait4Submit(future);
if (records.size() >= this.batchRecords) {
wait4Submit(future);
final List<Tuple2<Record, GdbElement>> batch = records;
future = submitService.submit(() -> batchCommitRecords(batch));
records = new ArrayList<>(batchRecords);
}
}
final List<Tuple2<Record, GdbElement>> batch = records;
future = this.submitService.submit(() -> batchCommitRecords(batch));
records = new ArrayList<>(this.batchRecords);
}
}
wait4Submit(future);
if (!records.isEmpty()) {
final List<Tuple2<Record, GdbElement>> batch = records;
future = submitService.submit(() -> batchCommitRecords(batch));
wait4Submit(future);
}
wait4Submit(future);
if (!records.isEmpty()) {
final List<Tuple2<Record, GdbElement>> batch = records;
future = this.submitService.submit(() -> batchCommitRecords(batch));
wait4Submit(future);
}
}
private void wait4Submit(Future<Boolean> future) {
if (future == null) {
return;
}
private void wait4Submit(final Future<Boolean> future) {
if (future == null) {
return;
}
try {
future.get();
} catch (Exception e) {
e.printStackTrace();
}
try {
future.get();
} catch (final Exception e) {
e.printStackTrace();
}
}
private boolean batchCommitRecords(final List<Tuple2<Record, GdbElement>> records) {
TaskPluginCollector collector = getTaskPluginCollector();
try {
List<Tuple2<Record, Exception>> errors = graph.add(records);
errors.forEach(t -> collector.collectDirtyRecord(t.getFirst(), t.getSecond()));
failed += errors.size();
} catch (Exception e) {
records.forEach(t -> collector.collectDirtyRecord(t.getFirst(), e));
failed += records.size();
}
final TaskPluginCollector collector = getTaskPluginCollector();
try {
final List<Tuple2<Record, Exception>> errors = this.graph.add(records);
errors.forEach(t -> collector.collectDirtyRecord(t.getFirst(), t.getSecond()));
this.failed += errors.size();
} catch (final Exception e) {
records.forEach(t -> collector.collectDirtyRecord(t.getFirst(), e));
this.failed += records.size();
}
records.clear();
return true;
records.clear();
return true;
}
@Override
@@ -231,7 +239,7 @@ public class GdbWriter extends Writer {
* 注意:此方法每个 Task 都会执行一次。
* 最佳实践:如果 Task 中有需要进行数据同步之后的后续处理,可以在此处完成。
*/
log.info("Task done, dirty record count - {}", failed);
log.info("Task done, dirty record count - {}", this.failed);
}
@Override
@@ -241,9 +249,9 @@ public class GdbWriter extends Writer {
* 最佳实践:通常配合Task 中的 post() 方法一起完成 Task 的资源释放。
*/
if (session) {
graph.close();
this.graph.close();
}
submitService.shutdown();
this.submitService.shutdown();
}
}
@@ -27,7 +27,6 @@ public enum GdbWriterErrorCode implements ErrorCode {
@Override
public String toString() {
return String.format("Code:[%s], Description:[%s]. ", this.code,
this.description);
return String.format("Code:[%s], Description:[%s]. ", this.code, this.description);
}
}
@@ -6,136 +6,164 @@ public final class Key {
* 此处声明插件用到的需要插件使用者提供的配置项
*/
public final static String HOST = "host";
public final static String PORT = "port";
public final static String HOST = "host";
public final static String PORT = "port";
public final static String USERNAME = "username";
public static final String PASSWORD = "password";
public static final String PASSWORD = "password";
/**
* import type and mode
*/
public static final String IMPORT_TYPE = "labelType";
public static final String UPDATE_MODE = "writeMode";
/**
* import type and mode
*/
public static final String IMPORT_TYPE = "labelType";
public static final String UPDATE_MODE = "writeMode";
/**
* label prefix issue
*/
public static final String ID_TRANS_RULE = "idTransRule";
public static final String SRC_ID_TRANS_RULE = "srcIdTransRule";
public static final String DST_ID_TRANS_RULE = "dstIdTransRule";
/**
* label prefix issue
*/
public static final String ID_TRANS_RULE = "idTransRule";
public static final String SRC_ID_TRANS_RULE = "srcIdTransRule";
public static final String DST_ID_TRANS_RULE = "dstIdTransRule";
public static final String LABEL = "label";
public static final String SRC_LABEL = "srcLabel";
public static final String DST_LABEL = "dstLabel";
public static final String LABEL = "label";
public static final String SRC_LABEL = "srcLabel";
public static final String DST_LABEL = "dstLabel";
public static final String MAPPING = "mapping";
public static final String MAPPING = "mapping";
/**
* column define in Gdb
*/
public static final String COLUMN = "column";
public static final String COLUMN_NAME = "name";
public static final String COLUMN_VALUE = "value";
public static final String COLUMN_TYPE = "type";
public static final String COLUMN_NODE_TYPE = "columnType";
/**
* column define in Gdb
*/
public static final String COLUMN = "column";
public static final String COLUMN_NAME = "name";
public static final String COLUMN_VALUE = "value";
public static final String COLUMN_TYPE = "type";
public static final String COLUMN_NODE_TYPE = "columnType";
/**
* Gdb Vertex/Edge elements
*/
public static final String ID = "id";
public static final String FROM = "from";
public static final String TO = "to";
public static final String PROPERTIES = "properties";
public static final String PROP_KEY = "name";
public static final String PROP_VALUE = "value";
public static final String PROP_TYPE = "type";
/**
* Gdb Vertex/Edge elements
*/
public static final String ID = "id";
public static final String FROM = "from";
public static final String TO = "to";
public static final String PROPERTIES = "properties";
public static final String PROP_KEY = "name";
public static final String PROP_VALUE = "value";
public static final String PROP_TYPE = "type";
public static final String PROPERTIES_JSON_STR = "propertiesJsonStr";
public static final String MAX_PROPERTIES_BATCH_NUM = "maxPropertiesBatchNumber";
public static final String PROPERTIES_JSON_STR = "propertiesJsonStr";
public static final String MAX_PROPERTIES_BATCH_NUM = "maxPropertiesBatchNumber";
/**
* session less client configure for connect pool
*/
public static final String MAX_IN_PROCESS_PER_CONNECTION = "maxInProcessPerConnection";
public static final String MAX_CONNECTION_POOL_SIZE = "maxConnectionPoolSize";
public static final String MAX_SIMULTANEOUS_USAGE_PER_CONNECTION = "maxSimultaneousUsagePerConnection";
/**
* session less client configure for connect pool
*/
public static final String MAX_IN_PROCESS_PER_CONNECTION = "maxInProcessPerConnection";
public static final String MAX_CONNECTION_POOL_SIZE = "maxConnectionPoolSize";
public static final String MAX_SIMULTANEOUS_USAGE_PER_CONNECTION = "maxSimultaneousUsagePerConnection";
public static final String MAX_RECORDS_IN_BATCH = "maxRecordsInBatch";
public static final String SESSION_STATE = "session";
public static final String MAX_RECORDS_IN_BATCH = "maxRecordsInBatch";
public static final String SESSION_STATE = "session";
public static enum ImportType {
/**
* Import vertices
*/
VERTEX,
/**
* Import edges
*/
EDGE;
}
public static enum UpdateMode {
/**
* Insert new records, fail if exists
*/
INSERT,
/**
* Skip this record if exists
*/
SKIP,
/**
* Update property of this record if exists
*/
MERGE;
}
/**
* request length limit, include gdb element string length GDB字段长度限制配置,可分别配置各字段的限制,超过限制的记录会当脏数据处理
*/
public static final String MAX_GDB_STRING_LENGTH = "maxStringLengthLimit";
public static final String MAX_GDB_ID_LENGTH = "maxIdStringLengthLimit";
public static final String MAX_GDB_LABEL_LENGTH = "maxLabelStringLengthLimit";
public static final String MAX_GDB_PROP_KEY_LENGTH = "maxPropKeyStringLengthLimit";
public static final String MAX_GDB_PROP_VALUE_LENGTH = "maxPropValueStringLengthLimit";
public static enum ColumnType {
/**
* vertex or edge id
*/
primaryKey,
public static final String MAX_GDB_REQUEST_LENGTH = "maxRequestLengthLimit";
/**
* vertex property
*/
vertexProperty,
public static enum ImportType {
/**
* Import vertices
*/
VERTEX,
/**
* Import edges
*/
EDGE;
}
/**
* start vertex id of edge
*/
srcPrimaryKey,
public static enum UpdateMode {
/**
* Insert new records, fail if exists
*/
INSERT,
/**
* Skip this record if exists
*/
SKIP,
/**
* Update property of this record if exists
*/
MERGE;
}
/**
* end vertex id of edge
*/
dstPrimaryKey,
public static enum ColumnType {
/**
* vertex or edge id
*/
primaryKey,
/**
* edge property
*/
edgeProperty,
/**
* vertex property
*/
vertexProperty,
/**
* vertex json style property
*/
vertexJsonProperty,
/**
* vertex setProperty
*/
vertexSetProperty,
/**
* edge json style property
*/
edgeJsonProperty
}
/**
* start vertex id of edge
*/
srcPrimaryKey,
public static enum IdTransRule {
/**
* vertex or edge id with 'label' prefix
*/
labelPrefix,
/**
* end vertex id of edge
*/
dstPrimaryKey,
/**
* vertex or edge id raw
*/
none
}
/**
* edge property
*/
edgeProperty,
/**
* vertex json style property
*/
vertexJsonProperty,
/**
* edge json style property
*/
edgeJsonProperty
}
public static enum IdTransRule {
/**
* vertex or edge id with 'label' prefix
*/
labelPrefix,
/**
* vertex or edge id raw
*/
none
}
public static enum PropertyType {
/**
* single Vertex Property
*/
single,
/**
* set Vertex Property
*/
set
}
}
@@ -3,37 +3,37 @@
*/
package com.alibaba.datax.plugin.writer.gdbwriter.client;
import java.util.ArrayList;
import java.util.List;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.writer.gdbwriter.model.GdbGraph;
import com.alibaba.datax.plugin.writer.gdbwriter.model.ScriptGdbGraph;
import java.util.ArrayList;
import java.util.List;
/**
* @author jerrywang
*
*/
public class GdbGraphManager implements AutoCloseable {
private static final GdbGraphManager instance = new GdbGraphManager();
private List<GdbGraph> graphs = new ArrayList<>();
public static GdbGraphManager instance() {
return instance;
}
private static final GdbGraphManager INSTANCE = new GdbGraphManager();
public GdbGraph getGraph(Configuration config, boolean session) {
GdbGraph graph = new ScriptGdbGraph(config, session);
graphs.add(graph);
return graph;
}
private List<GdbGraph> graphs = new ArrayList<>();
@Override
public void close() {
for(GdbGraph graph : graphs) {
graph.close();
}
graphs.clear();
}
public static GdbGraphManager instance() {
return INSTANCE;
}
public GdbGraph getGraph(final Configuration config, final boolean session) {
final GdbGraph graph = new ScriptGdbGraph(config, session);
this.graphs.add(graph);
return graph;
}
@Override
public void close() {
for (final GdbGraph graph : this.graphs) {
graph.close();
}
this.graphs.clear();
}
}
@@ -3,39 +3,43 @@
*/
package com.alibaba.datax.plugin.writer.gdbwriter.client;
import static com.alibaba.datax.plugin.writer.gdbwriter.util.ConfigHelper.assertConfig;
import static com.alibaba.datax.plugin.writer.gdbwriter.util.ConfigHelper.assertHasContent;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.writer.gdbwriter.Key;
import static com.alibaba.datax.plugin.writer.gdbwriter.util.ConfigHelper.*;
/**
* @author jerrywang
*
*/
public class GdbWriterConfig {
public static final int DEFAULT_MAX_IN_PROCESS_PER_CONNECTION = 4;
public static final int DEFAULT_MAX_CONNECTION_POOL_SIZE = 8;
public static final int DEFAULT_MAX_SIMULTANEOUS_USAGE_PER_CONNECTION = 8;
public static final int DEFAULT_BATCH_PROPERTY_NUM = 30;
public static final int DEFAULT_RECORD_NUM_IN_BATCH = 16;
public static final int DEFAULT_MAX_IN_PROCESS_PER_CONNECTION = 4;
public static final int DEFAULT_MAX_CONNECTION_POOL_SIZE = 8;
public static final int DEFAULT_MAX_SIMULTANEOUS_USAGE_PER_CONNECTION = 8;
public static final int DEFAULT_BATCH_PROPERTY_NUM = 30;
public static final int DEFAULT_RECORD_NUM_IN_BATCH = 16;
private Configuration config;
public static final int MAX_STRING_LENGTH = 10240;
public static final int MAX_REQUEST_LENGTH = 65535 - 1000;
private GdbWriterConfig(Configuration config) {
this.config = config;
private Configuration config;
validate();
}
private GdbWriterConfig(final Configuration config) {
this.config = config;
private void validate() {
assertHasContent(config, Key.HOST);
assertConfig(Key.PORT, () -> config.getInt(Key.PORT) > 0);
validate();
}
assertHasContent(config, Key.USERNAME);
assertHasContent(config, Key.PASSWORD);
}
public static GdbWriterConfig of(Configuration config) {
return new GdbWriterConfig(config);
}
public static GdbWriterConfig of(final Configuration config) {
return new GdbWriterConfig(config);
}
private void validate() {
assertHasContent(this.config, Key.HOST);
assertConfig(Key.PORT, () -> this.config.getInt(Key.PORT) > 0);
assertHasContent(this.config, Key.USERNAME);
assertHasContent(this.config, Key.PASSWORD);
}
}
@@ -3,6 +3,8 @@
*/
package com.alibaba.datax.plugin.writer.gdbwriter.mapping;
import static com.alibaba.datax.plugin.writer.gdbwriter.Key.ImportType.VERTEX;
import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
@@ -12,179 +14,191 @@ import java.util.regex.Matcher;
import java.util.regex.Pattern;
import com.alibaba.datax.common.element.Record;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.writer.gdbwriter.Key;
import com.alibaba.datax.plugin.writer.gdbwriter.model.GdbEdge;
import com.alibaba.datax.plugin.writer.gdbwriter.model.GdbElement;
import com.alibaba.datax.plugin.writer.gdbwriter.model.GdbVertex;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
import lombok.extern.slf4j.Slf4j;
import static com.alibaba.datax.plugin.writer.gdbwriter.Key.ImportType.VERTEX;
/**
* @author jerrywang
*
*/
@Slf4j
public class DefaultGdbMapper implements GdbMapper {
private static final Pattern STR_PATTERN = Pattern.compile("\\$\\{(\\d+)}");
private static final Pattern NORMAL_PATTERN = Pattern.compile("^\\$\\{(\\d+)}$");
private static final Pattern STR_DOLLAR_PATTERN = Pattern.compile("\\$\\{(\\d+)}");
private static final Pattern NORMAL_DOLLAR_PATTERN = Pattern.compile("^\\$\\{(\\d+)}$");
@Override
public Function<Record, GdbElement> getMapper(MappingRule rule) {
return r -> {
GdbElement e = (rule.getImportType() == VERTEX) ? new GdbVertex() : new GdbEdge();
forElement(rule).accept(r, e);
return e;
private static final Pattern STR_NUM_PATTERN = Pattern.compile("#\\{(\\d+)}");
private static final Pattern NORMAL_NUM_PATTERN = Pattern.compile("^#\\{(\\d+)}$");
public DefaultGdbMapper() {}
public DefaultGdbMapper(final Configuration config) {
MapperConfig.getInstance().updateConfig(config);
}
private static BiConsumer<Record, GdbElement> forElement(final MappingRule rule) {
final boolean numPattern = rule.isNumPattern();
final List<BiConsumer<Record, GdbElement>> properties = new ArrayList<>();
for (final MappingRule.PropertyMappingRule propRule : rule.getProperties()) {
final Function<Record, String> keyFunc = forStrColumn(numPattern, propRule.getKey());
if (propRule.getValueType() == ValueType.STRING) {
final Function<Record, String> valueFunc = forStrColumn(numPattern, propRule.getValue());
properties.add((r, e) -> {
e.addProperty(keyFunc.apply(r), valueFunc.apply(r), propRule.getPType());
});
} else {
final Function<Record, Object> valueFunc =
forObjColumn(numPattern, propRule.getValue(), propRule.getValueType());
properties.add((r, e) -> {
e.addProperty(keyFunc.apply(r), valueFunc.apply(r), propRule.getPType());
});
}
}
if (rule.getPropertiesJsonStr() != null) {
final Function<Record, String> jsonFunc = forStrColumn(numPattern, rule.getPropertiesJsonStr());
properties.add((r, e) -> {
final String propertiesStr = jsonFunc.apply(r);
final JSONObject root = (JSONObject)JSONObject.parse(propertiesStr);
final JSONArray propertiesList = root.getJSONArray("properties");
for (final Object object : propertiesList) {
final JSONObject jsonObject = (JSONObject)object;
final String key = jsonObject.getString("k");
final String name = jsonObject.getString("v");
final String type = jsonObject.getString("t");
final String card = jsonObject.getString("c");
if (key == null || name == null) {
continue;
}
addToProperties(e, key, name, type, card);
}
});
}
final BiConsumer<Record, GdbElement> ret = (r, e) -> {
final String label = forStrColumn(numPattern, rule.getLabel()).apply(r);
String id = forStrColumn(numPattern, rule.getId()).apply(r);
if (rule.getImportType() == Key.ImportType.EDGE) {
final String to = forStrColumn(numPattern, rule.getTo()).apply(r);
final String from = forStrColumn(numPattern, rule.getFrom()).apply(r);
if (to == null || from == null) {
log.error("invalid record to: {} , from: {}", to, from);
throw new IllegalArgumentException("to or from missed in edge");
}
((GdbEdge)e).setTo(to);
((GdbEdge)e).setFrom(from);
// generate UUID for edge
if (id == null) {
id = UUID.randomUUID().toString();
}
}
if (id == null || label == null) {
log.error("invalid record id: {} , label: {}", id, label);
throw new IllegalArgumentException("id or label missed");
}
e.setId(id);
e.setLabel(label);
properties.forEach(p -> p.accept(r, e));
};
}
return ret;
}
private static BiConsumer<Record, GdbElement> forElement(MappingRule rule) {
List<BiConsumer<Record, GdbElement>> properties = new ArrayList<>();
for (MappingRule.PropertyMappingRule propRule : rule.getProperties()) {
Function<Record, String> keyFunc = forStrColumn(propRule.getKey());
private static Function<Record, Object> forObjColumn(final boolean numPattern, final String rule, final ValueType type) {
final Pattern pattern = numPattern ? NORMAL_NUM_PATTERN : NORMAL_DOLLAR_PATTERN;
final Matcher m = pattern.matcher(rule);
if (m.matches()) {
final int index = Integer.valueOf(m.group(1));
return r -> type.applyColumn(r.getColumn(index));
} else {
return r -> type.fromStrFunc(rule);
}
}
if (propRule.getValueType() == ValueType.STRING) {
final Function<Record, String> valueFunc = forStrColumn(propRule.getValue());
properties.add((r, e) -> {
String k = keyFunc.apply(r);
String v = valueFunc.apply(r);
if (k != null && v != null) {
e.getProperties().put(k, v);
}
});
} else {
final Function<Record, Object> valueFunc = forObjColumn(propRule.getValue(), propRule.getValueType());
properties.add((r, e) -> {
String k = keyFunc.apply(r);
Object v = valueFunc.apply(r);
if (k != null && v != null) {
e.getProperties().put(k, v);
}
});
}
}
private static Function<Record, String> forStrColumn(final boolean numPattern, final String rule) {
final List<BiConsumer<StringBuilder, Record>> list = new ArrayList<>();
final Pattern pattern = numPattern ? STR_NUM_PATTERN : STR_DOLLAR_PATTERN;
final Matcher m = pattern.matcher(rule);
int last = 0;
while (m.find()) {
final String index = m.group(1);
// as simple integer index.
final int i = Integer.parseInt(index);
if (rule.getPropertiesJsonStr() != null) {
Function<Record, String> jsonFunc = forStrColumn(rule.getPropertiesJsonStr());
properties.add((r, e) -> {
String propertiesStr = jsonFunc.apply(r);
JSONObject root = (JSONObject)JSONObject.parse(propertiesStr);
JSONArray propertiesList = root.getJSONArray("properties");
final int tmp = last;
final int start = m.start();
list.add((sb, record) -> {
sb.append(rule.subSequence(tmp, start));
if (record.getColumn(i) != null && record.getColumn(i).getByteSize() > 0) {
sb.append(record.getColumn(i).asString());
}
});
for (Object object : propertiesList) {
JSONObject jsonObject = (JSONObject)object;
String key = jsonObject.getString("k");
String name = jsonObject.getString("v");
String type = jsonObject.getString("t");
last = m.end();
}
if (key == null || name == null) {
continue;
}
addToProperties(e, key, name, type);
}
});
}
final int tmp = last;
list.add((sb, record) -> {
sb.append(rule.subSequence(tmp, rule.length()));
});
BiConsumer<Record, GdbElement> ret = (r, e) -> {
String label = forStrColumn(rule.getLabel()).apply(r);
String id = forStrColumn(rule.getId()).apply(r);
return r -> {
final StringBuilder sb = new StringBuilder();
list.forEach(c -> c.accept(sb, r));
final String res = sb.toString();
return res.isEmpty() ? null : res;
};
}
if (rule.getImportType() == Key.ImportType.EDGE) {
String to = forStrColumn(rule.getTo()).apply(r);
String from = forStrColumn(rule.getFrom()).apply(r);
if (to == null || from == null) {
log.error("invalid record to: {} , from: {}", to, from);
throw new IllegalArgumentException("to or from missed in edge");
}
((GdbEdge)e).setTo(to);
((GdbEdge)e).setFrom(from);
private static boolean addToProperties(final GdbElement e, final String key, final String value, final String type, final String card) {
final Object pValue;
final ValueType valueType = ValueType.fromShortName(type);
// generate UUID for edge
if (id == null) {
id = UUID.randomUUID().toString();
}
}
if (valueType == ValueType.STRING) {
pValue = value;
} else if (valueType == ValueType.INT || valueType == ValueType.INTEGER) {
pValue = Integer.valueOf(value);
} else if (valueType == ValueType.LONG) {
pValue = Long.valueOf(value);
} else if (valueType == ValueType.DOUBLE) {
pValue = Double.valueOf(value);
} else if (valueType == ValueType.FLOAT) {
pValue = Float.valueOf(value);
} else if (valueType == ValueType.BOOLEAN) {
pValue = Boolean.valueOf(value);
} else {
log.error("invalid property key {}, value {}, type {}", key, value, type);
return false;
}
if (id == null || label == null) {
log.error("invalid record id: {} , label: {}", id, label);
throw new IllegalArgumentException("id or label missed");
}
// apply vertexSetProperty
if (Key.PropertyType.set.name().equals(card) && (e instanceof GdbVertex)) {
e.addProperty(key, pValue, Key.PropertyType.set);
} else {
e.addProperty(key, pValue);
}
return true;
}
e.setId(id);
e.setLabel(label);
properties.forEach(p -> p.accept(r, e));
};
return ret;
}
static Function<Record, Object> forObjColumn(String rule, ValueType type) {
Matcher m = NORMAL_PATTERN.matcher(rule);
if (m.matches()) {
int index = Integer.valueOf(m.group(1));
return r -> type.applyColumn(r.getColumn(index));
} else {
return r -> type.fromStrFunc(rule);
}
}
static Function<Record, String> forStrColumn(String rule) {
List<BiConsumer<StringBuilder, Record>> list = new ArrayList<>();
Matcher m = STR_PATTERN.matcher(rule);
int last = 0;
while (m.find()) {
String index = m.group(1);
// as simple integer index.
int i = Integer.parseInt(index);
final int tmp = last;
final int start = m.start();
list.add((sb, record) -> {
sb.append(rule.subSequence(tmp, start));
if(record.getColumn(i) != null && record.getColumn(i).getByteSize() > 0) {
sb.append(record.getColumn(i).asString());
}
});
last = m.end();
}
final int tmp = last;
list.add((sb, record) -> {
sb.append(rule.subSequence(tmp, rule.length()));
});
return r -> {
StringBuilder sb = new StringBuilder();
list.forEach(c -> c.accept(sb, r));
String res = sb.toString();
return res.isEmpty() ? null : res;
};
}
static boolean addToProperties(GdbElement e, String key, String value, String type) {
ValueType valueType = ValueType.fromShortName(type);
if(valueType == ValueType.STRING) {
e.getProperties().put(key, value);
} else if (valueType == ValueType.INT) {
e.getProperties().put(key, Integer.valueOf(value));
} else if (valueType == ValueType.LONG) {
e.getProperties().put(key, Long.valueOf(value));
} else if (valueType == ValueType.DOUBLE) {
e.getProperties().put(key, Double.valueOf(value));
} else if (valueType == ValueType.FLOAT) {
e.getProperties().put(key, Float.valueOf(value));
} else if (valueType == ValueType.BOOLEAN) {
e.getProperties().put(key, Boolean.valueOf(value));
} else {
log.error("invalid property key {}, value {}, type {}", key, value, type);
return false;
}
return true;
}
@Override
public Function<Record, GdbElement> getMapper(final MappingRule rule) {
return r -> {
final GdbElement e = (rule.getImportType() == VERTEX) ? new GdbVertex() : new GdbEdge();
forElement(rule).accept(r, e);
return e;
};
}
}
@@ -13,5 +13,5 @@ import com.alibaba.datax.plugin.writer.gdbwriter.model.GdbElement;
*
*/
public interface GdbMapper {
Function<Record, GdbElement> getMapper(MappingRule rule);
Function<Record, GdbElement> getMapper(MappingRule rule);
}
@@ -0,0 +1,68 @@
/*
* (C) 2019-present Alibaba Group Holding Limited.
*
* This program is free software; you can redistribute it and/or modify it under the terms of the GNU General Public
* License version 2 as published by the Free Software Foundation.
*/
package com.alibaba.datax.plugin.writer.gdbwriter.mapping;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.writer.gdbwriter.Key;
import com.alibaba.datax.plugin.writer.gdbwriter.client.GdbWriterConfig;
/**
* @author : Liu Jianping
* @date : 2019/10/15
*/
public class MapperConfig {
private static MapperConfig instance = new MapperConfig();
private int maxIdLength;
private int maxLabelLength;
private int maxPropKeyLength;
private int maxPropValueLength;
private MapperConfig() {
this.maxIdLength = GdbWriterConfig.MAX_STRING_LENGTH;
this.maxLabelLength = GdbWriterConfig.MAX_STRING_LENGTH;
this.maxPropKeyLength = GdbWriterConfig.MAX_STRING_LENGTH;
this.maxPropValueLength = GdbWriterConfig.MAX_STRING_LENGTH;
}
public static MapperConfig getInstance() {
return instance;
}
public void updateConfig(final Configuration config) {
final int length = config.getInt(Key.MAX_GDB_STRING_LENGTH, GdbWriterConfig.MAX_STRING_LENGTH);
Integer sLength = config.getInt(Key.MAX_GDB_ID_LENGTH);
this.maxIdLength = sLength == null ? length : sLength;
sLength = config.getInt(Key.MAX_GDB_LABEL_LENGTH);
this.maxLabelLength = sLength == null ? length : sLength;
sLength = config.getInt(Key.MAX_GDB_PROP_KEY_LENGTH);
this.maxPropKeyLength = sLength == null ? length : sLength;
sLength = config.getInt(Key.MAX_GDB_PROP_VALUE_LENGTH);
this.maxPropValueLength = sLength == null ? length : sLength;
}
public int getMaxIdLength() {
return this.maxIdLength;
}
public int getMaxLabelLength() {
return this.maxLabelLength;
}
public int getMaxPropKeyLength() {
return this.maxPropKeyLength;
}
public int getMaxPropValueLength() {
return this.maxPropValueLength;
}
}
@@ -7,6 +7,7 @@ import java.util.ArrayList;
import java.util.List;
import com.alibaba.datax.plugin.writer.gdbwriter.Key.ImportType;
import com.alibaba.datax.plugin.writer.gdbwriter.Key.PropertyType;
import lombok.Data;
@@ -16,26 +17,30 @@ import lombok.Data;
*/
@Data
public class MappingRule {
private String id = null;
private String id = null;
private String label = null;
private ImportType importType = null;
private String from = null;
private String label = null;
private String to = null;
private ImportType importType = null;
private List<PropertyMappingRule> properties = new ArrayList<>();
private String from = null;
private String propertiesJsonStr = null;
private String to = null;
@Data
public static class PropertyMappingRule {
private String key = null;
private String value = null;
private ValueType valueType = null;
}
private List<PropertyMappingRule> properties = new ArrayList<>();
private String propertiesJsonStr = null;
private boolean numPattern = false;
@Data
public static class PropertyMappingRule {
private String key = null;
private String value = null;
private ValueType valueType = null;
private PropertyType pType = PropertyType.single;
}
}
@@ -3,18 +3,21 @@
*/
package com.alibaba.datax.plugin.writer.gdbwriter.mapping;
import java.util.List;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import com.alibaba.datax.common.exception.DataXException;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.writer.gdbwriter.GdbWriterErrorCode;
import com.alibaba.datax.plugin.writer.gdbwriter.Key;
import com.alibaba.datax.plugin.writer.gdbwriter.Key.ImportType;
import com.alibaba.datax.plugin.writer.gdbwriter.Key.IdTransRule;
import com.alibaba.datax.plugin.writer.gdbwriter.Key.ColumnType;
import com.alibaba.datax.plugin.writer.gdbwriter.Key.IdTransRule;
import com.alibaba.datax.plugin.writer.gdbwriter.Key.ImportType;
import com.alibaba.datax.plugin.writer.gdbwriter.mapping.MappingRule.PropertyMappingRule;
import com.alibaba.datax.plugin.writer.gdbwriter.util.ConfigHelper;
import lombok.extern.slf4j.Slf4j;
import java.util.List;
import lombok.extern.slf4j.Slf4j;
/**
* @author jerrywang
@@ -22,66 +25,94 @@ import java.util.List;
*/
@Slf4j
public class MappingRuleFactory {
private static final MappingRuleFactory instance = new MappingRuleFactory();
public static final MappingRuleFactory getInstance() {
return instance;
}
private static final MappingRuleFactory instance = new MappingRuleFactory();
private static final Pattern STR_PATTERN = Pattern.compile("\\$\\{(\\d+)}");
private static final Pattern STR_NUM_PATTERN = Pattern.compile("#\\{(\\d+)}");
@Deprecated
public MappingRule create(Configuration config, ImportType type) {
MappingRule rule = new MappingRule();
rule.setId(config.getString(Key.ID));
rule.setLabel(config.getString(Key.LABEL));
if (type == ImportType.EDGE) {
rule.setFrom(config.getString(Key.FROM));
rule.setTo(config.getString(Key.TO));
}
rule.setImportType(type);
List<Configuration> configurations = config.getListConfiguration(Key.PROPERTIES);
if (configurations != null) {
for (Configuration prop : config.getListConfiguration(Key.PROPERTIES)) {
PropertyMappingRule propRule = new PropertyMappingRule();
propRule.setKey(prop.getString(Key.PROP_KEY));
propRule.setValue(prop.getString(Key.PROP_VALUE));
propRule.setValueType(ValueType.fromShortName(prop.getString(Key.PROP_TYPE).toLowerCase()));
rule.getProperties().add(propRule);
}
}
String propertiesJsonStr = config.getString(Key.PROPERTIES_JSON_STR, null);
if (propertiesJsonStr != null) {
rule.setPropertiesJsonStr(propertiesJsonStr);
}
return rule;
}
public MappingRule createV2(Configuration config) {
try {
ImportType type = ImportType.valueOf(config.getString(Key.IMPORT_TYPE));
return createV2(config, type);
} catch (NullPointerException e) {
throw DataXException.asDataXException(GdbWriterErrorCode.CONFIG_ITEM_MISS, Key.IMPORT_TYPE);
} catch (IllegalArgumentException e) {
throw DataXException.asDataXException(GdbWriterErrorCode.BAD_CONFIG_VALUE, Key.IMPORT_TYPE);
}
public static MappingRuleFactory getInstance() {
return instance;
}
public MappingRule createV2(Configuration config, ImportType type) {
MappingRule rule = new MappingRule();
private static boolean isPattern(final String value, final MappingRule rule, final boolean checked) {
if (checked) {
return true;
}
ConfigHelper.assertHasContent(config, Key.LABEL);
rule.setLabel(config.getString(Key.LABEL));
rule.setImportType(type);
if (value == null || value.isEmpty()) {
return false;
}
IdTransRule srcTransRule = IdTransRule.none;
Matcher m = STR_PATTERN.matcher(value);
if (m.find()) {
rule.setNumPattern(false);
return true;
}
m = STR_NUM_PATTERN.matcher(value);
if (m.find()) {
rule.setNumPattern(true);
return true;
}
return false;
}
@Deprecated
public MappingRule create(final Configuration config, final ImportType type) {
final MappingRule rule = new MappingRule();
rule.setId(config.getString(Key.ID));
rule.setLabel(config.getString(Key.LABEL));
if (type == ImportType.EDGE) {
rule.setFrom(config.getString(Key.FROM));
rule.setTo(config.getString(Key.TO));
}
rule.setImportType(type);
final List<Configuration> configurations = config.getListConfiguration(Key.PROPERTIES);
if (configurations != null) {
for (final Configuration prop : config.getListConfiguration(Key.PROPERTIES)) {
final PropertyMappingRule propRule = new PropertyMappingRule();
propRule.setKey(prop.getString(Key.PROP_KEY));
propRule.setValue(prop.getString(Key.PROP_VALUE));
propRule.setValueType(ValueType.fromShortName(prop.getString(Key.PROP_TYPE).toLowerCase()));
rule.getProperties().add(propRule);
}
}
final String propertiesJsonStr = config.getString(Key.PROPERTIES_JSON_STR, null);
if (propertiesJsonStr != null) {
rule.setPropertiesJsonStr(propertiesJsonStr);
}
return rule;
}
public MappingRule createV2(final Configuration config) {
try {
final ImportType type = ImportType.valueOf(config.getString(Key.IMPORT_TYPE));
return createV2(config, type);
} catch (final NullPointerException e) {
throw DataXException.asDataXException(GdbWriterErrorCode.CONFIG_ITEM_MISS, Key.IMPORT_TYPE);
} catch (final IllegalArgumentException e) {
throw DataXException.asDataXException(GdbWriterErrorCode.BAD_CONFIG_VALUE, Key.IMPORT_TYPE);
}
}
public MappingRule createV2(final Configuration config, final ImportType type) {
final MappingRule rule = new MappingRule();
boolean patternChecked = false;
ConfigHelper.assertHasContent(config, Key.LABEL);
rule.setLabel(config.getString(Key.LABEL));
rule.setImportType(type);
patternChecked = isPattern(rule.getLabel(), rule, patternChecked);
IdTransRule srcTransRule = IdTransRule.none;
IdTransRule dstTransRule = IdTransRule.none;
if (type == ImportType.EDGE) {
ConfigHelper.assertHasContent(config, Key.SRC_ID_TRANS_RULE);
ConfigHelper.assertHasContent(config, Key.DST_ID_TRANS_RULE);
ConfigHelper.assertHasContent(config, Key.SRC_ID_TRANS_RULE);
ConfigHelper.assertHasContent(config, Key.DST_ID_TRANS_RULE);
srcTransRule = IdTransRule.valueOf(config.getString(Key.SRC_ID_TRANS_RULE));
dstTransRule = IdTransRule.valueOf(config.getString(Key.DST_ID_TRANS_RULE));
@@ -94,88 +125,96 @@ public class MappingRuleFactory {
ConfigHelper.assertHasContent(config, Key.DST_LABEL);
}
}
ConfigHelper.assertHasContent(config, Key.ID_TRANS_RULE);
IdTransRule transRule = IdTransRule.valueOf(config.getString(Key.ID_TRANS_RULE));
ConfigHelper.assertHasContent(config, Key.ID_TRANS_RULE);
final IdTransRule transRule = IdTransRule.valueOf(config.getString(Key.ID_TRANS_RULE));
List<Configuration> configurationList = config.getListConfiguration(Key.COLUMN);
ConfigHelper.assertConfig(Key.COLUMN, () -> (configurationList != null && !configurationList.isEmpty()));
for (Configuration column : configurationList) {
ConfigHelper.assertHasContent(column, Key.COLUMN_NAME);
ConfigHelper.assertHasContent(column, Key.COLUMN_VALUE);
ConfigHelper.assertHasContent(column, Key.COLUMN_TYPE);
ConfigHelper.assertHasContent(column, Key.COLUMN_NODE_TYPE);
final List<Configuration> configurationList = config.getListConfiguration(Key.COLUMN);
ConfigHelper.assertConfig(Key.COLUMN, () -> (configurationList != null && !configurationList.isEmpty()));
for (final Configuration column : configurationList) {
ConfigHelper.assertHasContent(column, Key.COLUMN_NAME);
ConfigHelper.assertHasContent(column, Key.COLUMN_VALUE);
ConfigHelper.assertHasContent(column, Key.COLUMN_TYPE);
ConfigHelper.assertHasContent(column, Key.COLUMN_NODE_TYPE);
String columnValue = column.getString(Key.COLUMN_VALUE);
ColumnType columnType = ColumnType.valueOf(column.getString(Key.COLUMN_NODE_TYPE));
if (columnValue == null || columnValue.isEmpty()) {
// only allow edge empty id
ConfigHelper.assertConfig("empty column value",
() -> (type == ImportType.EDGE && columnType == ColumnType.primaryKey));
}
final String columnValue = column.getString(Key.COLUMN_VALUE);
final ColumnType columnType = ColumnType.valueOf(column.getString(Key.COLUMN_NODE_TYPE));
if (columnValue == null || columnValue.isEmpty()) {
// only allow edge empty id
ConfigHelper.assertConfig("empty column value",
() -> (type == ImportType.EDGE && columnType == ColumnType.primaryKey));
}
patternChecked = isPattern(columnValue, rule, patternChecked);
if (columnType == ColumnType.primaryKey) {
ValueType propType = ValueType.fromShortName(column.getString(Key.COLUMN_TYPE));
ConfigHelper.assertConfig("only string is allowed in primary key", () -> (propType == ValueType.STRING));
if (columnType == ColumnType.primaryKey) {
final ValueType propType = ValueType.fromShortName(column.getString(Key.COLUMN_TYPE));
ConfigHelper.assertConfig("only string is allowed in primary key",
() -> (propType == ValueType.STRING));
if (transRule == IdTransRule.labelPrefix) {
rule.setId(config.getString(Key.LABEL) + columnValue);
} else {
rule.setId(columnValue);
}
} else if (columnType == ColumnType.edgeJsonProperty || columnType == ColumnType.vertexJsonProperty) {
// only support one json property in column
ConfigHelper.assertConfig("multi JsonProperty", () -> (rule.getPropertiesJsonStr() == null));
rule.setPropertiesJsonStr(columnValue);
} else if (columnType == ColumnType.vertexProperty || columnType == ColumnType.edgeProperty) {
PropertyMappingRule propertyMappingRule = new PropertyMappingRule();
propertyMappingRule.setKey(column.getString(Key.COLUMN_NAME));
propertyMappingRule.setValue(columnValue);
ValueType propType = ValueType.fromShortName(column.getString(Key.COLUMN_TYPE));
ConfigHelper.assertConfig("unsupported property type", () -> propType != null);
propertyMappingRule.setValueType(propType);
rule.getProperties().add(propertyMappingRule);
} else if (columnType == ColumnType.srcPrimaryKey) {
if (type != ImportType.EDGE) {
continue;
}
ValueType propType = ValueType.fromShortName(column.getString(Key.COLUMN_TYPE));
ConfigHelper.assertConfig("only string is allowed in primary key", () -> (propType == ValueType.STRING));
if (srcTransRule == IdTransRule.labelPrefix) {
rule.setFrom(config.getString(Key.SRC_LABEL) + columnValue);
if (transRule == IdTransRule.labelPrefix) {
rule.setId(config.getString(Key.LABEL) + columnValue);
} else {
rule.setFrom(columnValue);
rule.setId(columnValue);
}
} else if (columnType == ColumnType.dstPrimaryKey) {
} else if (columnType == ColumnType.edgeJsonProperty || columnType == ColumnType.vertexJsonProperty) {
// only support one json property in column
ConfigHelper.assertConfig("multi JsonProperty", () -> (rule.getPropertiesJsonStr() == null));
rule.setPropertiesJsonStr(columnValue);
} else if (columnType == ColumnType.vertexProperty || columnType == ColumnType.edgeProperty
|| columnType == ColumnType.vertexSetProperty) {
final PropertyMappingRule propertyMappingRule = new PropertyMappingRule();
propertyMappingRule.setKey(column.getString(Key.COLUMN_NAME));
propertyMappingRule.setValue(columnValue);
final ValueType propType = ValueType.fromShortName(column.getString(Key.COLUMN_TYPE));
ConfigHelper.assertConfig("unsupported property type", () -> propType != null);
if (columnType == ColumnType.vertexSetProperty) {
propertyMappingRule.setPType(Key.PropertyType.set);
}
propertyMappingRule.setValueType(propType);
rule.getProperties().add(propertyMappingRule);
} else if (columnType == ColumnType.srcPrimaryKey) {
if (type != ImportType.EDGE) {
continue;
}
ValueType propType = ValueType.fromShortName(column.getString(Key.COLUMN_TYPE));
ConfigHelper.assertConfig("only string is allowed in primary key", () -> (propType == ValueType.STRING));
final ValueType propType = ValueType.fromShortName(column.getString(Key.COLUMN_TYPE));
ConfigHelper.assertConfig("only string is allowed in primary key",
() -> (propType == ValueType.STRING));
if (srcTransRule == IdTransRule.labelPrefix) {
rule.setFrom(config.getString(Key.SRC_LABEL) + columnValue);
} else {
rule.setFrom(columnValue);
}
} else if (columnType == ColumnType.dstPrimaryKey) {
if (type != ImportType.EDGE) {
continue;
}
final ValueType propType = ValueType.fromShortName(column.getString(Key.COLUMN_TYPE));
ConfigHelper.assertConfig("only string is allowed in primary key",
() -> (propType == ValueType.STRING));
if (dstTransRule == IdTransRule.labelPrefix) {
rule.setTo(config.getString(Key.DST_LABEL) + columnValue);
} else {
rule.setTo(columnValue);
}
}
}
}
}
if (rule.getImportType() == ImportType.EDGE) {
if (rule.getId() == null) {
rule.setId("");
log.info("edge id is missed, uuid be default");
}
ConfigHelper.assertConfig("to needed in edge", () -> (rule.getTo() != null));
ConfigHelper.assertConfig("from needed in edge", () -> (rule.getFrom() != null));
}
ConfigHelper.assertConfig("id needed", () -> (rule.getId() != null));
if (rule.getImportType() == ImportType.EDGE) {
if (rule.getId() == null) {
rule.setId("");
log.info("edge id is missed, uuid be default");
}
ConfigHelper.assertConfig("to needed in edge", () -> (rule.getTo() != null));
ConfigHelper.assertConfig("from needed in edge", () -> (rule.getFrom() != null));
}
ConfigHelper.assertConfig("id needed", () -> (rule.getId() != null));
return rule;
}
return rule;
}
}
@@ -8,6 +8,7 @@ import java.util.Map;
import java.util.function.Function;
import com.alibaba.datax.common.element.Column;
import lombok.extern.slf4j.Slf4j;
/**
@@ -16,56 +17,61 @@ import lombok.extern.slf4j.Slf4j;
*/
@Slf4j
public enum ValueType {
INT(Integer.class, "int", Column::asLong, Integer::valueOf),
LONG(Long.class, "long", Column::asLong, Long::valueOf),
DOUBLE(Double.class, "double", Column::asDouble, Double::valueOf),
FLOAT(Float.class, "float", Column::asDouble, Float::valueOf),
BOOLEAN(Boolean.class, "boolean", Column::asBoolean, Boolean::valueOf),
STRING(String.class, "string", Column::asString, String::valueOf);
/**
* property value type
*/
INT(Integer.class, "int", Column::asLong, Integer::valueOf),
INTEGER(Integer.class, "integer", Column::asLong, Integer::valueOf),
LONG(Long.class, "long", Column::asLong, Long::valueOf),
DOUBLE(Double.class, "double", Column::asDouble, Double::valueOf),
FLOAT(Float.class, "float", Column::asDouble, Float::valueOf),
BOOLEAN(Boolean.class, "boolean", Column::asBoolean, Boolean::valueOf),
STRING(String.class, "string", Column::asString, String::valueOf);
private Class<?> type = null;
private String shortName = null;
private Function<Column, Object> columnFunc = null;
private Function<String, Object> fromStrFunc = null;
private Class<?> type = null;
private String shortName = null;
private Function<Column, Object> columnFunc = null;
private Function<String, Object> fromStrFunc = null;
private ValueType(Class<?> type, String name, Function<Column, Object> columnFunc, Function<String, Object> fromStrFunc) {
this.type = type;
this.shortName = name;
this.columnFunc = columnFunc;
this.fromStrFunc = fromStrFunc;
ValueTypeHolder.shortName2type.put(name, this);
}
public static ValueType fromShortName(String name) {
return ValueTypeHolder.shortName2type.get(name);
}
private ValueType(final Class<?> type, final String name, final Function<Column, Object> columnFunc,
final Function<String, Object> fromStrFunc) {
this.type = type;
this.shortName = name;
this.columnFunc = columnFunc;
this.fromStrFunc = fromStrFunc;
public Class<?> type() {
return this.type;
}
public String shortName() {
return this.shortName;
}
public Object applyColumn(Column column) {
try {
if (column == null) {
return null;
}
return columnFunc.apply(column);
} catch (Exception e) {
log.error("applyColumn error {}, column {}", e.toString(), column);
throw e;
}
}
public Object fromStrFunc(String str) {
return fromStrFunc.apply(str);
}
ValueTypeHolder.shortName2type.put(name, this);
}
private static class ValueTypeHolder {
private static Map<String, ValueType> shortName2type = new HashMap<>();
}
public static ValueType fromShortName(final String name) {
return ValueTypeHolder.shortName2type.get(name);
}
public Class<?> type() {
return this.type;
}
public String shortName() {
return this.shortName;
}
public Object applyColumn(final Column column) {
try {
if (column == null) {
return null;
}
return this.columnFunc.apply(column);
} catch (final Exception e) {
log.error("applyColumn error {}, column {}", e.toString(), column);
throw e;
}
}
public Object fromStrFunc(final String str) {
return this.fromStrFunc.apply(str);
}
private static class ValueTypeHolder {
private static Map<String, ValueType> shortName2type = new HashMap<>();
}
}
@@ -3,20 +3,24 @@
*/
package com.alibaba.datax.plugin.writer.gdbwriter.model;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.writer.gdbwriter.Key;
import com.alibaba.datax.plugin.writer.gdbwriter.client.GdbWriterConfig;
import static com.alibaba.datax.plugin.writer.gdbwriter.client.GdbWriterConfig.DEFAULT_BATCH_PROPERTY_NUM;
import static com.alibaba.datax.plugin.writer.gdbwriter.client.GdbWriterConfig.MAX_REQUEST_LENGTH;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import lombok.extern.slf4j.Slf4j;
import org.apache.tinkerpop.gremlin.driver.Client;
import org.apache.tinkerpop.gremlin.driver.Cluster;
import org.apache.tinkerpop.gremlin.driver.RequestOptions;
import org.apache.tinkerpop.gremlin.driver.ResultSet;
import org.apache.tinkerpop.gremlin.driver.ser.Serializers;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.writer.gdbwriter.Key;
import com.alibaba.datax.plugin.writer.gdbwriter.client.GdbWriterConfig;
import lombok.extern.slf4j.Slf4j;
/**
* @author jerrywang
@@ -24,128 +28,124 @@ import java.util.concurrent.TimeUnit;
*/
@Slf4j
public abstract class AbstractGdbGraph implements GdbGraph {
private final static int DEFAULT_TIMEOUT = 30000;
private final static int DEFAULT_TIMEOUT = 30000;
protected Client client = null;
protected Key.UpdateMode updateMode = Key.UpdateMode.INSERT;
protected int propertiesBatchNum = GdbWriterConfig.DEFAULT_BATCH_PROPERTY_NUM;
protected boolean session = false;
protected Client client = null;
protected Key.UpdateMode updateMode = Key.UpdateMode.INSERT;
protected int propertiesBatchNum = DEFAULT_BATCH_PROPERTY_NUM;
protected boolean session = false;
protected int maxRequestLength = GdbWriterConfig.MAX_REQUEST_LENGTH;
protected AbstractGdbGraph() {}
protected AbstractGdbGraph() {}
protected AbstractGdbGraph(final Configuration config, final boolean session) {
initClient(config, session);
}
protected AbstractGdbGraph(Configuration config, boolean session) {
initClient(config, session);
}
protected void initClient(final Configuration config, final boolean session) {
this.updateMode = Key.UpdateMode.valueOf(config.getString(Key.UPDATE_MODE, "INSERT"));
log.info("init graphdb client");
final String host = config.getString(Key.HOST);
final int port = config.getInt(Key.PORT);
final String username = config.getString(Key.USERNAME);
final String password = config.getString(Key.PASSWORD);
int maxDepthPerConnection =
config.getInt(Key.MAX_IN_PROCESS_PER_CONNECTION, GdbWriterConfig.DEFAULT_MAX_IN_PROCESS_PER_CONNECTION);
protected void initClient(Configuration config, boolean session) {
updateMode = Key.UpdateMode.valueOf(config.getString(Key.UPDATE_MODE, "INSERT"));
log.info("init graphdb client");
String host = config.getString(Key.HOST);
int port = config.getInt(Key.PORT);
String username = config.getString(Key.USERNAME);
String password = config.getString(Key.PASSWORD);
int maxDepthPerConnection = config.getInt(Key.MAX_IN_PROCESS_PER_CONNECTION,
GdbWriterConfig.DEFAULT_MAX_IN_PROCESS_PER_CONNECTION);
int maxConnectionPoolSize =
config.getInt(Key.MAX_CONNECTION_POOL_SIZE, GdbWriterConfig.DEFAULT_MAX_CONNECTION_POOL_SIZE);
int maxConnectionPoolSize = config.getInt(Key.MAX_CONNECTION_POOL_SIZE,
GdbWriterConfig.DEFAULT_MAX_CONNECTION_POOL_SIZE);
int maxSimultaneousUsagePerConnection = config.getInt(Key.MAX_SIMULTANEOUS_USAGE_PER_CONNECTION,
GdbWriterConfig.DEFAULT_MAX_SIMULTANEOUS_USAGE_PER_CONNECTION);
int maxSimultaneousUsagePerConnection = config.getInt(Key.MAX_SIMULTANEOUS_USAGE_PER_CONNECTION,
GdbWriterConfig.DEFAULT_MAX_SIMULTANEOUS_USAGE_PER_CONNECTION);
this.session = session;
if (this.session) {
maxConnectionPoolSize = GdbWriterConfig.DEFAULT_MAX_CONNECTION_POOL_SIZE;
maxDepthPerConnection = GdbWriterConfig.DEFAULT_MAX_IN_PROCESS_PER_CONNECTION;
maxSimultaneousUsagePerConnection = GdbWriterConfig.DEFAULT_MAX_SIMULTANEOUS_USAGE_PER_CONNECTION;
}
this.session = session;
if (this.session) {
maxConnectionPoolSize = GdbWriterConfig.DEFAULT_MAX_CONNECTION_POOL_SIZE;
maxDepthPerConnection = GdbWriterConfig.DEFAULT_MAX_IN_PROCESS_PER_CONNECTION;
maxSimultaneousUsagePerConnection = GdbWriterConfig.DEFAULT_MAX_SIMULTANEOUS_USAGE_PER_CONNECTION;
}
try {
final Cluster cluster = Cluster.build(host).port(port).credentials(username, password)
.serializer(Serializers.GRAPHBINARY_V1D0).maxContentLength(1048576)
.maxInProcessPerConnection(maxDepthPerConnection).minInProcessPerConnection(0)
.maxConnectionPoolSize(maxConnectionPoolSize).minConnectionPoolSize(maxConnectionPoolSize)
.maxSimultaneousUsagePerConnection(maxSimultaneousUsagePerConnection).resultIterationBatchSize(64)
.create();
this.client = session ? cluster.connect(UUID.randomUUID().toString()).init() : cluster.connect().init();
warmClient(maxConnectionPoolSize * maxDepthPerConnection);
} catch (final RuntimeException e) {
log.error("Failed to connect to GDB {}:{}, due to {}", host, port, e);
throw e;
}
try {
Cluster cluster = Cluster.build(host).port(port).credentials(username, password)
.serializer(Serializers.GRAPHBINARY_V1D0)
.maxContentLength(1048576)
.maxInProcessPerConnection(maxDepthPerConnection)
.minInProcessPerConnection(0)
.maxConnectionPoolSize(maxConnectionPoolSize)
.minConnectionPoolSize(maxConnectionPoolSize)
.maxSimultaneousUsagePerConnection(maxSimultaneousUsagePerConnection)
.resultIterationBatchSize(64)
.create();
client = session ? cluster.connect(UUID.randomUUID().toString()).init() : cluster.connect().init();
warmClient(maxConnectionPoolSize*maxDepthPerConnection);
} catch (RuntimeException e) {
log.error("Failed to connect to GDB {}:{}, due to {}", host, port, e);
throw e;
}
this.propertiesBatchNum = config.getInt(Key.MAX_PROPERTIES_BATCH_NUM, DEFAULT_BATCH_PROPERTY_NUM);
this.maxRequestLength = config.getInt(Key.MAX_GDB_REQUEST_LENGTH, MAX_REQUEST_LENGTH);
}
propertiesBatchNum = config.getInt(Key.MAX_PROPERTIES_BATCH_NUM, GdbWriterConfig.DEFAULT_BATCH_PROPERTY_NUM);
}
/**
* @param dsl
* @param parameters
*/
protected void runInternal(final String dsl, final Map<String, Object> parameters) throws Exception {
final RequestOptions.Builder options = RequestOptions.build().timeout(DEFAULT_TIMEOUT);
if (parameters != null && !parameters.isEmpty()) {
parameters.forEach(options::addParameter);
}
final ResultSet results = this.client.submitAsync(dsl, options.create()).get(DEFAULT_TIMEOUT, TimeUnit.MILLISECONDS);
results.all().get(DEFAULT_TIMEOUT + 1000, TimeUnit.MILLISECONDS);
}
/**
* @param dsl
* @param parameters
*/
protected void runInternal(String dsl, final Map<String, Object> parameters) throws Exception {
RequestOptions.Builder options = RequestOptions.build().timeout(DEFAULT_TIMEOUT);
if (parameters != null && !parameters.isEmpty()) {
parameters.forEach(options::addParameter);
}
void beginTx() {
if (!this.session) {
return;
}
ResultSet results = client.submitAsync(dsl, options.create()).get(DEFAULT_TIMEOUT, TimeUnit.MILLISECONDS);
results.all().get(DEFAULT_TIMEOUT + 1000, TimeUnit.MILLISECONDS);
}
final String dsl = "g.tx().open()";
this.client.submit(dsl).all().join();
}
void beginTx() {
if (!session) {
return;
}
void doCommit() {
if (!this.session) {
return;
}
String dsl = "g.tx().open()";
client.submit(dsl).all().join();
}
try {
final String dsl = "g.tx().commit()";
this.client.submit(dsl).all().join();
} catch (final Exception e) {
throw new RuntimeException(e);
}
}
void doCommit() {
if (!session) {
return;
}
void doRollback() {
if (!this.session) {
return;
}
try {
String dsl = "g.tx().commit()";
client.submit(dsl).all().join();
} catch (Exception e) {
throw new RuntimeException(e);
}
}
final String dsl = "g.tx().rollback()";
this.client.submit(dsl).all().join();
}
void doRollback() {
if (!session) {
return;
}
private void warmClient(final int num) {
try {
beginTx();
runInternal("g.V('test')", null);
doCommit();
log.info("warm graphdb client over");
} catch (final Exception e) {
log.error("warmClient error");
doRollback();
throw new RuntimeException(e);
}
}
String dsl = "g.tx().rollback()";
client.submit(dsl).all().join();
}
private void warmClient(int num) {
try {
beginTx();
runInternal("g.V('test')", null);
doCommit();
log.info("warm graphdb client over");
} catch (Exception e) {
log.error("warmClient error");
doRollback();
throw new RuntimeException(e);
}
}
@Override
public void close() {
if (client != null) {
log.info("close graphdb client");
client.close();
}
}
@Override
public void close() {
if (this.client != null) {
log.info("close graphdb client");
this.client.close();
}
}
}
@@ -3,7 +3,8 @@
*/
package com.alibaba.datax.plugin.writer.gdbwriter.model;
import lombok.Data;
import com.alibaba.datax.plugin.writer.gdbwriter.mapping.MapperConfig;
import lombok.EqualsAndHashCode;
import lombok.ToString;
@@ -11,10 +12,33 @@ import lombok.ToString;
* @author jerrywang
*
*/
@Data
@EqualsAndHashCode(callSuper = true)
@ToString(callSuper = true)
public class GdbEdge extends GdbElement {
private String from = null;
private String to = null;
private String from = null;
private String to = null;
public String getFrom() {
return this.from;
}
public void setFrom(final String from) {
final int maxIdLength = MapperConfig.getInstance().getMaxIdLength();
if (from.length() > maxIdLength) {
throw new IllegalArgumentException("from length over limit(" + maxIdLength + ")");
}
this.from = from;
}
public String getTo() {
return this.to;
}
public void setTo(final String to) {
final int maxIdLength = MapperConfig.getInstance().getMaxIdLength();
if (to.length() > maxIdLength) {
throw new IllegalArgumentException("to length over limit(" + maxIdLength + ")");
}
this.to = to;
}
}
@@ -3,18 +3,107 @@
*/
package com.alibaba.datax.plugin.writer.gdbwriter.model;
import java.util.HashMap;
import java.util.Map;
import java.util.LinkedList;
import java.util.List;
import lombok.Data;
import com.alibaba.datax.plugin.writer.gdbwriter.Key.PropertyType;
import com.alibaba.datax.plugin.writer.gdbwriter.mapping.MapperConfig;
/**
* @author jerrywang
*
*/
@Data
public class GdbElement {
String id = null;
String label = null;
Map<String, Object> properties = new HashMap<>();
private String id = null;
private String label = null;
private List<GdbProperty> properties = new LinkedList<>();
public String getId() {
return this.id;
}
public void setId(final String id) {
final int maxIdLength = MapperConfig.getInstance().getMaxIdLength();
if (id.length() > maxIdLength) {
throw new IllegalArgumentException("id length over limit(" + maxIdLength + ")");
}
this.id = id;
}
public String getLabel() {
return this.label;
}
public void setLabel(final String label) {
final int maxLabelLength = MapperConfig.getInstance().getMaxLabelLength();
if (label.length() > maxLabelLength) {
throw new IllegalArgumentException("label length over limit(" + maxLabelLength + ")");
}
this.label = label;
}
public List<GdbProperty> getProperties() {
return this.properties;
}
public void addProperty(final String propKey, final Object propValue, final PropertyType card) {
if (propKey == null || propValue == null) {
return;
}
final int maxPropKeyLength = MapperConfig.getInstance().getMaxPropKeyLength();
if (propKey.length() > maxPropKeyLength) {
throw new IllegalArgumentException("property key length over limit(" + maxPropKeyLength + ")");
}
if (propValue instanceof String) {
final int maxPropValueLength = MapperConfig.getInstance().getMaxPropValueLength();
if (((String)propValue).length() > maxPropKeyLength) {
throw new IllegalArgumentException("property value length over limit(" + maxPropValueLength + ")");
}
}
this.properties.add(new GdbProperty(propKey, propValue, card));
}
public void addProperty(final String propKey, final Object propValue) {
addProperty(propKey, propValue, PropertyType.single);
}
@Override
public String toString() {
final StringBuffer sb = new StringBuffer(this.id + "[" + this.label + "]{");
this.properties.forEach(n -> {
sb.append(n.cardinality.name());
sb.append("[");
sb.append(n.key);
sb.append(" - ");
sb.append(String.valueOf(n.value));
sb.append("]");
});
return sb.toString();
}
public static class GdbProperty {
private String key;
private Object value;
private PropertyType cardinality;
private GdbProperty(final String key, final Object value, final PropertyType card) {
this.key = key;
this.value = value;
this.cardinality = card;
}
public PropertyType getCardinality() {
return this.cardinality;
}
public String getKey() {
return this.key;
}
public Object getValue() {
return this.value;
}
}
}
@@ -3,18 +3,19 @@
*/
package com.alibaba.datax.plugin.writer.gdbwriter.model;
import com.alibaba.datax.common.element.Record;
import groovy.lang.Tuple2;
import java.util.List;
import com.alibaba.datax.common.element.Record;
import groovy.lang.Tuple2;
/**
* @author jerrywang
*
*/
public interface GdbGraph extends AutoCloseable {
List<Tuple2<Record, Exception>> add(List<Tuple2<Record, GdbElement>> records);
List<Tuple2<Record, Exception>> add(List<Tuple2<Record, GdbElement>> records);
@Override
void close();
@Override
void close();
}
@@ -3,15 +3,17 @@
*/
package com.alibaba.datax.plugin.writer.gdbwriter.model;
import java.util.*;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Random;
import com.alibaba.datax.common.element.Record;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.writer.gdbwriter.Key;
import com.alibaba.datax.plugin.writer.gdbwriter.util.GdbDuplicateIdException;
import com.github.benmanes.caffeine.cache.Cache;
import com.github.benmanes.caffeine.cache.Caffeine;
import groovy.lang.Tuple2;
import lombok.extern.slf4j.Slf4j;
@@ -21,176 +23,198 @@ import lombok.extern.slf4j.Slf4j;
*/
@Slf4j
public class ScriptGdbGraph extends AbstractGdbGraph {
private static final String VAR_PREFIX = "GDB___";
private static final String VAR_ID = VAR_PREFIX + "id";
private static final String VAR_LABEL = VAR_PREFIX + "label";
private static final String VAR_FROM = VAR_PREFIX + "from";
private static final String VAR_TO = VAR_PREFIX + "to";
private static final String VAR_PROP_KEY = VAR_PREFIX + "PK";
private static final String VAR_PROP_VALUE = VAR_PREFIX + "PV";
private static final String ADD_V_START = "g.addV(" + VAR_LABEL + ").property(id, " + VAR_ID + ")";
private static final String ADD_E_START = "g.addE(" + VAR_LABEL + ").property(id, " + VAR_ID + ").from(V("
+ VAR_FROM + ")).to(V(" + VAR_TO + "))";
private static final String VAR_PREFIX = "GDB___";
private static final String VAR_ID = VAR_PREFIX + "id";
private static final String VAR_LABEL = VAR_PREFIX + "label";
private static final String VAR_FROM = VAR_PREFIX + "from";
private static final String VAR_TO = VAR_PREFIX + "to";
private static final String VAR_PROP_KEY = VAR_PREFIX + "PK";
private static final String VAR_PROP_VALUE = VAR_PREFIX + "PV";
private static final String ADD_V_START = "g.addV(" + VAR_LABEL + ").property(id, " + VAR_ID + ")";
private static final String ADD_E_START =
"g.addE(" + VAR_LABEL + ").property(id, " + VAR_ID + ").from(V(" + VAR_FROM + ")).to(V(" + VAR_TO + "))";
private static final String UPDATE_V_START = "g.V("+VAR_ID+")";
private static final String UPDATE_E_START = "g.E("+VAR_ID+")";
private static final String UPDATE_V_START = "g.V(" + VAR_ID + ")";
private static final String UPDATE_E_START = "g.E(" + VAR_ID + ")";
private Cache<Integer, String> propertyCache;
private Random random;
private Random random;
public ScriptGdbGraph() {
propertyCache = Caffeine.newBuilder().maximumSize(1024).build();
random = new Random();
}
public ScriptGdbGraph() {
this.random = new Random();
}
public ScriptGdbGraph(Configuration config, boolean session) {
super(config, session);
public ScriptGdbGraph(final Configuration config, final boolean session) {
super(config, session);
propertyCache = Caffeine.newBuilder().maximumSize(1024).build();
random = new Random();
this.random = new Random();
log.info("Init as ScriptGdbGraph.");
}
log.info("Init as ScriptGdbGraph.");
}
/**
* Apply list of {@link GdbElement} to GDB, return the failed records
*
* @param records
* list of element to apply
* @return
*/
@Override
public List<Tuple2<Record, Exception>> add(final List<Tuple2<Record, GdbElement>> records) {
final List<Tuple2<Record, Exception>> errors = new ArrayList<>();
try {
beginTx();
for (final Tuple2<Record, GdbElement> elementTuple2 : records) {
try {
addInternal(elementTuple2.getSecond());
} catch (final Exception e) {
errors.add(new Tuple2<>(elementTuple2.getFirst(), e));
}
}
doCommit();
} catch (final Exception ex) {
doRollback();
throw new RuntimeException(ex);
}
return errors;
}
/**
* Apply list of {@link GdbElement} to GDB, return the failed records
* @param records list of element to apply
* @return
*/
@Override
public List<Tuple2<Record, Exception>> add(List<Tuple2<Record, GdbElement>> records) {
List<Tuple2<Record, Exception>> errors = new ArrayList<>();
try {
beginTx();
for (Tuple2<Record, GdbElement> elementTuple2 : records) {
try {
addInternal(elementTuple2.getSecond());
} catch (Exception e) {
errors.add(new Tuple2<>(elementTuple2.getFirst(), e));
}
}
doCommit();
} catch (Exception ex) {
doRollback();
throw new RuntimeException(ex);
}
return errors;
}
private void addInternal(final GdbElement element) {
try {
addInternal(element, false);
} catch (final GdbDuplicateIdException e) {
if (this.updateMode == Key.UpdateMode.SKIP) {
log.debug("Skip duplicate id {}", element.getId());
} else if (this.updateMode == Key.UpdateMode.INSERT) {
throw new RuntimeException(e);
} else if (this.updateMode == Key.UpdateMode.MERGE) {
if (element.getProperties().isEmpty()) {
return;
}
private void addInternal(GdbElement element) {
try {
addInternal(element, false);
} catch (GdbDuplicateIdException e) {
if (updateMode == Key.UpdateMode.SKIP) {
log.debug("Skip duplicate id {}", element.getId());
} else if (updateMode == Key.UpdateMode.INSERT) {
throw new RuntimeException(e);
} else if (updateMode == Key.UpdateMode.MERGE) {
if (element.getProperties().isEmpty()) {
return;
}
try {
addInternal(element, true);
} catch (final GdbDuplicateIdException e1) {
log.error("duplicate id {} while update...", element.getId());
throw new RuntimeException(e1);
}
}
}
}
try {
addInternal(element, true);
} catch (GdbDuplicateIdException e1) {
log.error("duplicate id {} while update...", element.getId());
throw new RuntimeException(e1);
}
}
}
}
private void addInternal(final GdbElement element, final boolean update) throws GdbDuplicateIdException {
boolean firstAdd = !update;
final boolean isVertex = (element instanceof GdbVertex);
final List<GdbElement.GdbProperty> params = element.getProperties();
final List<GdbElement.GdbProperty> subParams = new ArrayList<>(this.propertiesBatchNum);
private void addInternal(GdbElement element, boolean update) throws GdbDuplicateIdException {
Map<String, Object> params = element.getProperties();
Map<String, Object> subParams = new HashMap<>(propertiesBatchNum);
boolean firstAdd = !update;
boolean isVertex = (element instanceof GdbVertex);
final int idLength = element.getId().length();
int attachLength = element.getLabel().length();
if (element instanceof GdbEdge) {
attachLength += ((GdbEdge)element).getFrom().length();
attachLength += ((GdbEdge)element).getTo().length();
}
for (Map.Entry<String, Object> entry : params.entrySet()) {
subParams.put(entry.getKey(), entry.getValue());
if (subParams.size() >= propertiesBatchNum) {
setGraphDbElement(element, subParams, isVertex, firstAdd);
firstAdd = false;
subParams.clear();
}
}
if (!subParams.isEmpty() || firstAdd) {
setGraphDbElement(element, subParams, isVertex, firstAdd);
}
}
int requestLength = idLength;
for (final GdbElement.GdbProperty entry : params) {
final String propKey = entry.getKey();
final Object propValue = entry.getValue();
private Tuple2<String, Map<String, Object>> buildDsl(GdbElement element,
Map<String, Object> properties,
boolean isVertex, boolean firstAdd) {
Map<String, Object> params = new HashMap<>();
int appendLength = propKey.length();
if (propValue instanceof String) {
appendLength += ((String)propValue).length();
}
String dslPropertyPart = propertyCache.get(properties.size(), keys -> {
final StringBuilder sb = new StringBuilder();
for (int i = 0; i < keys; i++) {
sb.append(".property(").append(VAR_PROP_KEY).append(i)
.append(", ").append(VAR_PROP_VALUE).append(i).append(")");
}
return sb.toString();
});
if (checkSplitDsl(firstAdd, requestLength, attachLength, appendLength, subParams.size())) {
setGraphDbElement(element, subParams, isVertex, firstAdd);
firstAdd = false;
subParams.clear();
requestLength = idLength;
}
String dsl;
if (isVertex) {
dsl = (firstAdd ? ADD_V_START : UPDATE_V_START) + dslPropertyPart;
} else {
dsl = (firstAdd ? ADD_E_START : UPDATE_E_START) + dslPropertyPart;
if (firstAdd) {
params.put(VAR_FROM, ((GdbEdge)element).getFrom());
params.put(VAR_TO, ((GdbEdge)element).getTo());
}
}
requestLength += appendLength;
subParams.add(entry);
}
if (!subParams.isEmpty() || firstAdd) {
checkSplitDsl(firstAdd, requestLength, attachLength, 0, 0);
setGraphDbElement(element, subParams, isVertex, firstAdd);
}
}
int index = 0;
for (Map.Entry<String, Object> entry : properties.entrySet()) {
params.put(VAR_PROP_KEY+index, entry.getKey());
params.put(VAR_PROP_VALUE+index, entry.getValue());
index++;
}
private boolean checkSplitDsl(final boolean firstAdd, final int requestLength, final int attachLength, final int appendLength,
final int propNum) {
final int length = firstAdd ? requestLength + attachLength : requestLength;
if (length > this.maxRequestLength) {
throw new IllegalArgumentException("request length over limit(" + this.maxRequestLength + ")");
}
return length + appendLength > this.maxRequestLength || propNum >= this.propertiesBatchNum;
}
if (firstAdd) {
params.put(VAR_LABEL, element.getLabel());
}
params.put(VAR_ID, element.getId());
private Tuple2<String, Map<String, Object>> buildDsl(final GdbElement element, final List<GdbElement.GdbProperty> properties,
final boolean isVertex, final boolean firstAdd) {
final Map<String, Object> params = new HashMap<>();
final StringBuilder sb = new StringBuilder();
if (isVertex) {
sb.append(firstAdd ? ADD_V_START : UPDATE_V_START);
} else {
sb.append(firstAdd ? ADD_E_START : UPDATE_E_START);
}
return new Tuple2<>(dsl, params);
}
for (int i = 0; i < properties.size(); i++) {
final GdbElement.GdbProperty prop = properties.get(i);
private void setGraphDbElement(GdbElement element, Map<String, Object> properties,
boolean isVertex, boolean firstAdd) throws GdbDuplicateIdException {
int retry = 10;
int idleTime = random.nextInt(10) + 10;
Tuple2<String, Map<String, Object>> elementDsl = buildDsl(element, properties, isVertex, firstAdd);
sb.append(".property(");
if (prop.getCardinality() == Key.PropertyType.set) {
sb.append("set, ");
}
sb.append(VAR_PROP_KEY).append(i).append(", ").append(VAR_PROP_VALUE).append(i).append(")");
while (retry > 0) {
try {
runInternal(elementDsl.getFirst(), elementDsl.getSecond());
log.debug("AddElement {}", element.getId());
return;
} catch (Exception e) {
String cause = e.getCause() == null ? "" : e.getCause().toString();
if (cause.contains("rejected from")) {
retry--;
try {
Thread.sleep(idleTime);
} catch (InterruptedException e1) {
// ...
}
idleTime = Math.min(idleTime * 2, 2000);
continue;
} else if (firstAdd && cause.contains("GraphDB id exists")) {
throw new GdbDuplicateIdException(e);
}
log.error("Add Failed id {}, dsl {}, params {}, e {}", element.getId(),
elementDsl.getFirst(), elementDsl.getSecond(), e);
throw new RuntimeException(e);
}
}
log.error("Add Failed id {}, dsl {}, params {}", element.getId(),
elementDsl.getFirst(), elementDsl.getSecond());
throw new RuntimeException("failed to queue new element to server");
}
params.put(VAR_PROP_KEY + i, prop.getKey());
params.put(VAR_PROP_VALUE + i, prop.getValue());
}
if (firstAdd) {
params.put(VAR_LABEL, element.getLabel());
if (!isVertex) {
params.put(VAR_FROM, ((GdbEdge)element).getFrom());
params.put(VAR_TO, ((GdbEdge)element).getTo());
}
}
params.put(VAR_ID, element.getId());
return new Tuple2<>(sb.toString(), params);
}
private void setGraphDbElement(final GdbElement element, final List<GdbElement.GdbProperty> properties, final boolean isVertex,
final boolean firstAdd) throws GdbDuplicateIdException {
int retry = 10;
int idleTime = this.random.nextInt(10) + 10;
final Tuple2<String, Map<String, Object>> elementDsl = buildDsl(element, properties, isVertex, firstAdd);
while (retry > 0) {
try {
runInternal(elementDsl.getFirst(), elementDsl.getSecond());
log.debug("AddElement {}", element.getId());
return;
} catch (final Exception e) {
final String cause = e.getCause() == null ? "" : e.getCause().toString();
if (cause.contains("rejected from") || cause.contains("Timeout waiting to lock key")) {
retry--;
try {
Thread.sleep(idleTime);
} catch (final InterruptedException e1) {
// ...
}
idleTime = Math.min(idleTime * 2, 2000);
continue;
} else if (firstAdd && cause.contains("GraphDB id exists")) {
throw new GdbDuplicateIdException(e);
}
log.error("Add Failed id {}, dsl {}, params {}, e {}", element.getId(), elementDsl.getFirst(),
elementDsl.getSecond(), e);
throw new RuntimeException(e);
}
}
log.error("Add Failed id {}, dsl {}, params {}", element.getId(), elementDsl.getFirst(),
elementDsl.getSecond());
throw new RuntimeException("failed to queue new element to server");
}
}
@@ -7,53 +7,57 @@ import java.io.IOException;
import java.io.InputStream;
import java.util.function.Supplier;
import org.apache.commons.lang3.StringUtils;
import com.alibaba.datax.common.exception.DataXException;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.writer.gdbwriter.GdbWriterErrorCode;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import org.apache.commons.lang3.StringUtils;
/**
* @author jerrywang
*
*/
public interface ConfigHelper {
static void assertConfig(String key, Supplier<Boolean> f) {
if (!f.get()) {
throw DataXException.asDataXException(GdbWriterErrorCode.BAD_CONFIG_VALUE, key);
}
}
static void assertConfig(final String key, final Supplier<Boolean> f) {
if (!f.get()) {
throw DataXException.asDataXException(GdbWriterErrorCode.BAD_CONFIG_VALUE, key);
}
}
static void assertHasContent(Configuration config, String key) {
assertConfig(key, () -> StringUtils.isNotBlank(config.getString(key)));
}
static void assertHasContent(final Configuration config, final String key) {
assertConfig(key, () -> StringUtils.isNotBlank(config.getString(key)));
}
/**
* NOTE: {@code Configuration::get(String, Class<T>)} doesn't work.
*
* @param conf Configuration
* @param key key path to configuration
* @param cls Class of result type
* @return the target configuration object of type T
*/
static <T> T getConfig(Configuration conf, String key, Class<T> cls) {
JSONObject j = (JSONObject) conf.get(key);
return JSON.toJavaObject(j, cls);
}
/**
* Create a configuration from the specified file on the classpath.
*
* @param name file name
* @return Configuration instance.
*/
static Configuration fromClasspath(String name) {
try (InputStream is = Thread.currentThread().getContextClassLoader().getResourceAsStream(name)) {
return Configuration.from(is);
} catch (IOException e) {
throw new IllegalArgumentException("File not found: " + name);
}
}
/**
* NOTE: {@code Configuration::get(String, Class<T>)} doesn't work.
*
* @param conf
* Configuration
* @param key
* key path to configuration
* @param cls
* Class of result type
* @return the target configuration object of type T
*/
static <T> T getConfig(final Configuration conf, final String key, final Class<T> cls) {
final JSONObject j = (JSONObject)conf.get(key);
return JSON.toJavaObject(j, cls);
}
/**
* Create a configuration from the specified file on the classpath.
*
* @param name
* file name
* @return Configuration instance.
*/
static Configuration fromClasspath(final String name) {
try (final InputStream is = Thread.currentThread().getContextClassLoader().getResourceAsStream(name)) {
return Configuration.from(is);
} catch (final IOException e) {
throw new IllegalArgumentException("File not found: " + name);
}
}
}
@@ -1,9 +1,8 @@
/*
* (C) 2019-present Alibaba Group Holding Limited.
* (C) 2019-present Alibaba Group Holding Limited.
*
* This program is free software; you can redistribute it and/or modify
* it under the terms of the GNU General Public License version 2 as
* published by the Free Software Foundation.
* This program is free software; you can redistribute it and/or modify it under the terms of the GNU General Public
* License version 2 as published by the Free Software Foundation.
*/
package com.alibaba.datax.plugin.writer.gdbwriter.util;
@@ -13,11 +12,11 @@ package com.alibaba.datax.plugin.writer.gdbwriter.util;
*/
public class GdbDuplicateIdException extends Exception {
public GdbDuplicateIdException(Exception e) {
super(e);
}
public GdbDuplicateIdException(Exception e) {
super(e);
}
public GdbDuplicateIdException() {
super();
}
public GdbDuplicateIdException() {
super();
}
}
Binary file not shown.

After

Width:  |  Height:  |  Size: 193 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 195 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 189 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 191 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 189 KiB

+1
View File
@@ -132,6 +132,7 @@ MongoDBReader通过Datax框架从MongoDB并行的读取数据,通过主控的J
* nameColumn的名字。【必填】
* typeColumn的类型。【选填】
* splitter:因为MongoDB支持数组类型,但是Datax框架本身不支持数组类型,所以mongoDB读出来的数组类型要通过这个分隔符合并成字符串。【选填】
* query: MongoDB的额外查询条件。【选填】
#### 5 类型转换
+1 -1
View File
@@ -147,7 +147,7 @@ MysqlWriter 通过 DataX 框架获取 Reader 生成的协议数据,根据你
* **column**
* 描述:目的表需要写入数据的字段,字段之间用英文逗号分隔。例如: "column": ["id","name","age"]。如果要依次写入全部列,使用*表示, 例如: "column": ["*"]。
* 描述:目的表需要写入数据的字段,字段之间用英文逗号分隔。例如: "column": ["id","name","age"]。如果要依次写入全部列,使用`*`表示, 例如: `"column": ["*"]`
**column配置项必须指定,不能留空!**
+2 -1
View File
@@ -58,7 +58,8 @@ ODPSReader 支持读取分区表、非分区表,不支持读取虚拟视图。
],
"packageAuthorizedProject": "yourCurrentProjectName",
"splitMode": "record",
"odpsServer": "http://xxx/api"
"odpsServer": "http://xxx/api",
"tunnelServer": "http://dt.odps.aliyun.com"
}
},
"writer": {
+5 -5
View File
@@ -43,11 +43,11 @@
<scope>system</scope>
<systemPath>${basedir}/src/main/libs/bcprov-jdk15on-1.52.jar</systemPath>
</dependency>
<dependency>
<groupId>com.aliyun.odps</groupId>
<artifactId>odps-sdk-core</artifactId>
<version>0.19.3-public</version>
</dependency>
<dependency>
<groupId>com.aliyun.odps</groupId>
<artifactId>odps-sdk-core</artifactId>
<version>0.20.7-public</version>
</dependency>
<dependency>
<groupId>org.mockito</groupId>
+44 -40
View File
@@ -33,47 +33,51 @@ ODPSWriter插件用于实现往ODPS插入或者更新数据,主要提供给etl
```json
{
"job": {
"setting": {
"speed": {"byte": 1048576}
"job": {
"setting": {
"speed": {
"byte": 1048576
}
},
"content": [
{
"reader": {
"name": "streamreader",
"parameter": {
"column": [
{
"value": "DataX",
"type": "string"
},
{
"value": "test",
"type": "bytes"
}
],
"sliceRecordCount": 100000
}
},
"content": [
{
"reader": {
"name": "streamreader",
"parameter": {
"column" : [
{
"value": "DataX",
"type": "string"
},
{
"value": "test",
"type": "bytes"
}
],
"sliceRecordCount": 100000
}
},
"writer": {
"name": "odpswriter",
"parameter": {
"project": "chinan_test",
"table": "odps_write_test00_partitioned",
"partition":"school=SiChuan-School,class=1",
"column": ["id","name"],
"accessId": "xxx",
"accessKey": "xxxx",
"truncate": true,
"odpsServer": "http://sxxx/api",
"tunnelServer": "http://xxx",
"accountType": "aliyun"
}
}
}
}
]
}
"writer": {
"name": "odpswriter",
"parameter": {
"project": "chinan_test",
"table": "odps_write_test00_partitioned",
"partition": "school=SiChuan-School,class=1",
"column": [
"id",
"name"
],
"accessId": "xxx",
"accessKey": "xxxx",
"truncate": true,
"odpsServer": "http://sxxx/api",
"tunnelServer": "http://xxx",
"accountType": "aliyun"
}
}
}
]
}
}
```
+1 -1
View File
@@ -40,7 +40,7 @@
<dependency>
<groupId>com.aliyun.odps</groupId>
<artifactId>odps-sdk-core</artifactId>
<version>0.19.3-public</version>
<version>0.20.7-public</version>
</dependency>
<!-- httpclient begin -->
+1 -1
View File
@@ -101,7 +101,7 @@ OracleReader插件实现了从Oracle读取数据。在底层实现上,OracleRe
"connection": [
{
"querySql": [
"select db_id,on_line_flag from db_info where db_id < 10;"
"select db_id,on_line_flag from db_info where db_id < 10"
],
"jdbcUrl": [
"jdbc:oracle:thin:@[HOST_NAME]:PORT:[DATABASE_NAME]"
+3 -2
View File
@@ -48,12 +48,13 @@ OSSWriter实现了从DataX协议转为OSS中的TXT文件功能,OSS本身是无
},
"writer": {
"name": "osswriter",
"parameter": {
"endpoint": "http://oss.aliyuncs.com",
"accessId": "",
"accessKey": "",
"bucket": "myBucket",
"object": "/cdo/datax",
"object": "cdo/datax",
"encoding": "UTF-8",
"fieldDelimiter": ",",
"writeMode": "truncate|append|nonConflict"
@@ -104,7 +105,7 @@ OSSWriter实现了从DataX协议转为OSS中的TXT文件功能,OSS本身是无
* 描述:OSSWriter写入的文件名,OSS使用文件名模拟目录的实现。 <br />
使用"object": "datax",写入object以datax开头,后缀添加随机字符串。
使用"object": "/cdo/datax",写入的object以/cdo/datax开头,后缀随机添加字符串,/作为OSS模拟目录的分隔符。
使用"object": "cdo/datax",写入的object以cdo/datax开头,后缀随机添加字符串,/作为OSS模拟目录的分隔符。
* 必选:是 <br />
+2 -2
View File
@@ -10,7 +10,7 @@
</parent>
<groupId>com.alibaba.datax</groupId>
<artifactId>otsstreamreader</artifactId>
<version>0.0.1-SNAPSHOT</version>
<version>0.0.1</version>
<dependencies>
<dependency>
@@ -32,7 +32,7 @@
<dependency>
<groupId>com.aliyun.openservices</groupId>
<artifactId>tablestore-streamclient</artifactId>
<version>1.0.0-SNAPSHOT</version>
<version>1.0.0</version>
</dependency>
<dependency>
<groupId>com.google.code.gson</groupId>
+7
View File
@@ -357,5 +357,12 @@
</includes>
<outputDirectory>datax</outputDirectory>
</fileSet>
<fileSet>
<directory>clickhousewriter/target/datax/</directory>
<includes>
<include>**/*.*</include>
</includes>
<outputDirectory>datax</outputDirectory>
</fileSet>
</fileSets>
</assembly>
@@ -18,7 +18,8 @@ public enum DataBaseType {
PostgreSQL("postgresql", "org.postgresql.Driver"),
RDBMS("rdbms", "com.alibaba.datax.plugin.rdbms.util.DataBaseType"),
DB2("db2", "com.ibm.db2.jcc.DB2Driver"),
ADS("ads","com.mysql.jdbc.Driver");
ADS("ads","com.mysql.jdbc.Driver"),
ClickHouse("clickhouse", "ru.yandex.clickhouse.ClickHouseDriver");
private String typeName;
@@ -54,6 +55,8 @@ public enum DataBaseType {
break;
case PostgreSQL:
break;
case ClickHouse:
break;
case RDBMS:
break;
default:
@@ -91,6 +94,8 @@ public enum DataBaseType {
break;
case PostgreSQL:
break;
case ClickHouse:
break;
case RDBMS:
break;
default:
+32 -2
View File
@@ -22,7 +22,7 @@
<commons-lang3-version>3.3.2</commons-lang3-version>
<commons-configuration-version>1.10</commons-configuration-version>
<commons-cli-version>1.2</commons-cli-version>
<fastjson-version>1.1.46.sec01</fastjson-version>
<fastjson-version>1.1.46.sec10</fastjson-version>
<guava-version>16.0.1</guava-version>
<diamond.version>3.7.2.1-SNAPSHOT</diamond.version>
@@ -63,8 +63,10 @@
<module>rdbmsreader</module>
<module>hbase11xreader</module>
<module>hbase094xreader</module>
<module>tsdbreader</module>
<module>opentsdbreader</module>
<module>cassandrareader</module>
<module>gdbreader</module>
<!-- writer -->
<module>mysqlwriter</module>
@@ -92,7 +94,7 @@
<module>adbpgwriter</module>
<module>gdbwriter</module>
<module>cassandrawriter</module>
<module>clickhousewriter</module>
<!-- common support module -->
<module>plugin-rdbms-util</module>
<module>plugin-unstructured-storage-util</module>
@@ -176,6 +178,34 @@
</dependencies>
</dependencyManagement>
<repositories>
<repository>
<id>central</id>
<name>Nexus aliyun</name>
<url>https://maven.aliyun.com/repository/central</url>
<releases>
<enabled>true</enabled>
</releases>
<snapshots>
<enabled>true</enabled>
</snapshots>
</repository>
</repositories>
<pluginRepositories>
<pluginRepository>
<id>central</id>
<name>Nexus aliyun</name>
<url>https://maven.aliyun.com/repository/central</url>
<releases>
<enabled>true</enabled>
</releases>
<snapshots>
<enabled>true</enabled>
</snapshots>
</pluginRepository>
</pluginRepositories>
<build>
<plugins>
<plugin>
+1 -1
View File
@@ -141,7 +141,7 @@ PostgresqlWriter通过 DataX 框架获取 Reader 生成的协议数据,根据
* **column**
* 描述:目的表需要写入数据的字段,字段之间用英文逗号分隔。例如: "column": ["id","name","age"]。如果要依次写入全部列,使用*表示, 例如: "column": ["*"]
* 描述:目的表需要写入数据的字段,字段之间用英文逗号分隔。例如: "column": ["id","name","age"]。如果要依次写入全部列,使用\*表示, 例如: "column": ["\*"]
注意:1、我们强烈不推荐你这样配置,因为当你目的表字段个数、类型等有改动时,你的任务可能运行不正确或者失败
2、此处 column 不能配置任何常量值
@@ -5,6 +5,7 @@ import com.alibaba.datax.common.plugin.RecordSender;
import com.alibaba.datax.common.spi.Reader;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.rdbms.reader.CommonRdbmsReader;
import com.alibaba.datax.plugin.rdbms.util.DBUtil;
import com.alibaba.datax.plugin.rdbms.util.DBUtilErrorCode;
import com.alibaba.datax.plugin.rdbms.util.DataBaseType;
@@ -12,7 +13,10 @@ import java.util.List;
public class RdbmsReader extends Reader {
private static final DataBaseType DATABASE_TYPE = DataBaseType.RDBMS;
static {
//加载插件下面配置的驱动类
DBUtil.loadDriverClass("reader", "rdbms");
}
public static class Job extends Reader.Job {
private Configuration originalConfig;
@@ -4,6 +4,7 @@ import com.alibaba.datax.common.exception.DataXException;
import com.alibaba.datax.common.plugin.RecordReceiver;
import com.alibaba.datax.common.spi.Writer;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.rdbms.util.DBUtil;
import com.alibaba.datax.plugin.rdbms.util.DBUtilErrorCode;
import com.alibaba.datax.plugin.rdbms.util.DataBaseType;
import com.alibaba.datax.plugin.rdbms.writer.CommonRdbmsWriter;
@@ -13,7 +14,10 @@ import java.util.List;
public class RdbmsWriter extends Writer {
private static final DataBaseType DATABASE_TYPE = DataBaseType.RDBMS;
static {
//加载插件下面配置的驱动类
DBUtil.loadDriverClass("writer", "rdbms");
}
public static class Job extends Writer.Job {
private Configuration originalConfig = null;
private CommonRdbmsWriter.Job commonRdbmsWriterMaster;
+1 -1
View File
@@ -2,7 +2,7 @@
## Transformer定义
在数据同步、传输过程中,存在用户对于数据传输进行特殊定制化的需求场景,包括裁剪列、转换列等工作,可以借助ETL的T过程实现(Transformer)。DataX包含了完的E(Extract)、T(Transformer)、L(Load)支持。
在数据同步、传输过程中,存在用户对于数据传输进行特殊定制化的需求场景,包括裁剪列、转换列等工作,可以借助ETL的T过程实现(Transformer)。DataX包含了完的E(Extract)、T(Transformer)、L(Load)支持。
## 运行模型
+587
View File
@@ -0,0 +1,587 @@
# TSDBReader 插件文档
___
## 1 快速介绍
TSDBReader 插件实现了从阿里云 TSDB 读取数据。阿里云时间序列数据库 ( **T**ime **S**eries **D**ata**b**ase , 简称 TSDB) 是一种集时序数据高效读写,压缩存储,实时计算能力为一体的数据库服务,可广泛应用于物联网和互联网领域,实现对设备及业务服务的实时监控,实时预测告警。详见 TSDB 的阿里云[官网](https://cn.aliyun.com/product/hitsdb)。
## 2 实现原理
在底层实现上,TSDBReader 通过 HTTP 请求链接到 阿里云 TSDB 实例,利用 `/api/query` 或者 `/api/mquery` 接口将数据点扫描出来(更多细节详见:[时序数据库 TSDB - HTTP API 概览](https://help.aliyun.com/document_detail/63557.html))。而整个同步的过程,是通过时间线和查询时间线范围进行切分。
## 3 功能说明
### 3.1 配置样例
* 配置一个从 阿里云 TSDB 数据库同步抽取数据到本地的作业,并以**时序数据**的格式输出:
时序数据样例:
```json
{"metric":"m","tags":{"app":"a19","cluster":"c5","group":"g10","ip":"i999","zone":"z1"},"timestamp":1546272263,"value":1}
```
```json
{
"job": {
"content": [
{
"reader": {
"name": "tsdbreader",
"parameter": {
"sinkDbType": "TSDB",
"endpoint": "http://localhost:8242",
"column": [
"m"
],
"splitIntervalMs": 60000,
"beginDateTime": "2019-01-01 00:00:00",
"endDateTime": "2019-01-01 01:00:00"
}
},
"writer": {
"name": "streamwriter",
"parameter": {
"encoding": "UTF-8",
"print": true
}
}
}
],
"setting": {
"speed": {
"channel": 3
}
}
}
}
```
* 配置一个从 阿里云 TSDB 数据库同步抽取数据到本地的作业,并以**关系型数据**的格式输出:
关系型数据样例:
```txt
m 1546272125 a1 c1 g2 i3021 z4 1.0
```
```json
{
"job": {
"content": [
{
"reader": {
"name": "tsdbreader",
"parameter": {
"sinkDbType": "RDB",
"endpoint": "http://localhost:8242",
"column": [
"__metric__",
"__ts__",
"app",
"cluster",
"group",
"ip",
"zone",
"__value__"
],
"metric": [
"m"
],
"splitIntervalMs": 60000,
"beginDateTime": "2019-01-01 00:00:00",
"endDateTime": "2019-01-01 01:00:00"
}
},
"writer": {
"name": "streamwriter",
"parameter": {
"encoding": "UTF-8",
"print": true
}
}
}
],
"setting": {
"speed": {
"channel": 3
}
}
}
}
```
* 配置一个从 阿里云 TSDB 数据库同步抽取**单值**数据到 ADB 的作业:
```json
{
"job": {
"content": [
{
"reader": {
"name": "tsdbreader",
"parameter": {
"sinkDbType": "RDB",
"endpoint": "http://localhost:8242",
"column": [
"__metric__",
"__ts__",
"app",
"cluster",
"group",
"ip",
"zone",
"__value__"
],
"metric": [
"m"
],
"splitIntervalMs": 60000,
"beginDateTime": "2019-01-01 00:00:00",
"endDateTime": "2019-01-01 01:00:00"
}
},
"writer": {
"name": "adswriter",
"parameter": {
"username": "******",
"password": "******",
"column": [
"`metric`",
"`ts`",
"`app`",
"`cluster`",
"`group`",
"`ip`",
"`zone`",
"`value`"
],
"url": "http://localhost:3306",
"schema": "datax_test",
"table": "datax_test",
"writeMode": "insert",
"opIndex": "0",
"batchSize": "2"
}
}
}
],
"setting": {
"speed": {
"channel": 3
}
}
}
}
```
* 配置一个从 阿里云 TSDB 数据库同步抽取**多值**数据到 ADB 的作业:
```json
{
"job": {
"content": [
{
"reader": {
"name": "tsdbreader",
"parameter": {
"sinkDbType": "RDB",
"endpoint": "http://localhost:8242",
"column": [
"__metric__",
"__ts__",
"app",
"cluster",
"group",
"ip",
"zone",
"load",
"memory",
"cpu"
],
"metric": [
"m_field"
],
"field": {
"m_field": [
"load",
"memory",
"cpu"
]
},
"splitIntervalMs": 60000,
"beginDateTime": "2019-01-01 00:00:00",
"endDateTime": "2019-01-01 01:00:00"
}
},
"writer": {
"name": "adswriter",
"parameter": {
"username": "******",
"password": "******",
"column": [
"`metric`",
"`ts`",
"`app`",
"`cluster`",
"`group`",
"`ip`",
"`zone`",
"`load`",
"`memory`",
"`cpu`"
],
"url": "http://localhost:3306",
"schema": "datax_test",
"table": "datax_test_multi_field",
"writeMode": "insert",
"opIndex": "0",
"batchSize": "2"
}
}
}
],
"setting": {
"speed": {
"channel": 3
}
}
}
}
```
* 配置一个从 阿里云 TSDB 数据库同步抽取**单值**数据到 ADB 的作业,并指定过滤部分时间线:
```json
{
"job": {
"content": [
{
"reader": {
"name": "tsdbreader",
"parameter": {
"sinkDbType": "RDB",
"endpoint": "http://localhost:8242",
"column": [
"__metric__",
"__ts__",
"app",
"cluster",
"group",
"ip",
"zone",
"__value__"
],
"metric": [
"m"
],
"tag": {
"m": {
"app": "a1",
"cluster": "c1"
}
},
"splitIntervalMs": 60000,
"beginDateTime": "2019-01-01 00:00:00",
"endDateTime": "2019-01-01 01:00:00"
}
},
"writer": {
"name": "adswriter",
"parameter": {
"username": "******",
"password": "******",
"column": [
"`metric`",
"`ts`",
"`app`",
"`cluster`",
"`group`",
"`ip`",
"`zone`",
"`value`"
],
"url": "http://localhost:3306",
"schema": "datax_test",
"table": "datax_test",
"writeMode": "insert",
"opIndex": "0",
"batchSize": "2"
}
}
}
],
"setting": {
"speed": {
"channel": 3
}
}
}
}
```
* 配置一个从 阿里云 TSDB 数据库同步抽取**多值**数据到 ADB 的作业,并指定过滤部分时间线:
```json
{
"job": {
"content": [
{
"reader": {
"name": "tsdbreader",
"parameter": {
"sinkDbType": "RDB",
"endpoint": "http://localhost:8242",
"column": [
"__metric__",
"__ts__",
"app",
"cluster",
"group",
"ip",
"zone",
"load",
"memory",
"cpu"
],
"metric": [
"m_field"
],
"field": {
"m_field": [
"load",
"memory",
"cpu"
]
},
"tag": {
"m_field": {
"ip": "i999"
}
},
"splitIntervalMs": 60000,
"beginDateTime": "2019-01-01 00:00:00",
"endDateTime": "2019-01-01 01:00:00"
}
},
"writer": {
"name": "adswriter",
"parameter": {
"username": "******",
"password": "******",
"column": [
"`metric`",
"`ts`",
"`app`",
"`cluster`",
"`group`",
"`ip`",
"`zone`",
"`load`",
"`memory`",
"`cpu`"
],
"url": "http://localhost:3306",
"schema": "datax_test",
"table": "datax_test_multi_field",
"writeMode": "insert",
"opIndex": "0",
"batchSize": "2"
}
}
}
],
"setting": {
"speed": {
"channel": 3
}
}
}
}
```
* 配置一个从 阿里云 TSDB 数据库同步抽取**单值**数据到另一个 阿里云 TSDB 数据库 的作业:
```json
{
"job": {
"content": [
{
"reader": {
"name": "tsdbreader",
"parameter": {
"sinkDbType": "TSDB",
"endpoint": "http://localhost:8242",
"column": [
"m"
],
"splitIntervalMs": 60000,
"beginDateTime": "2019-01-01 00:00:00",
"endDateTime": "2019-01-01 01:00:00"
}
},
"writer": {
"name": "tsdbwriter",
"parameter": {
"endpoint": "http://localhost:8240"
}
}
}
],
"setting": {
"speed": {
"channel": 3
}
}
}
}
```
* 配置一个从 阿里云 TSDB 数据库同步抽取**多值**数据到另一个 阿里云 TSDB 数据库 的作业:
```json
{
"job": {
"content": [
{
"reader": {
"name": "tsdbreader",
"parameter": {
"sinkDbType": "TSDB",
"endpoint": "http://localhost:8242",
"column": [
"m_field"
],
"field": {
"m_field": [
"load",
"memory",
"cpu"
]
},
"splitIntervalMs": 60000,
"beginDateTime": "2019-01-01 00:00:00",
"endDateTime": "2019-01-01 01:00:00"
}
},
"writer": {
"name": "tsdbwriter",
"parameter": {
"multiField": true,
"endpoint": "http://localhost:8240"
}
}
}
],
"setting": {
"speed": {
"channel": 3
}
}
}
}
```
### 3.2 参数说明
* **name**
* 描述:本插件的名称
* 必选:是
* 默认值:tsdbreader
* **parameter**
* **sinkDbType**
* 描述:目标数据库的类型
* 必选:否
* 默认值:TSDB
* 注意:目前支持 TSDB 和 RDB 两个取值。其中,TSDB 包括 阿里云 TSDB、OpenTSDB、InfluxDB、Prometheus 和 TimeScale。RDB 包括 ADB、MySQL、Oracle、PostgreSQL 和 DRDS 等。
* **endpoint**
* 描述:阿里云 TSDB 的 HTTP 连接地址
* 必选:是
* 格式:http://IP:Port
* 默认值:无
* **column**
* 描述:TSDB 场景下:数据迁移任务需要迁移的 Metric 列表;RDB 场景下:映射到关系型数据库中的表字段,且增加 `__metric__``__ts__``__value__` 三个字段,其中 `__metric__` 用于映射度量字段,`__ts__` 用于映射 timestamp 字段,而 `__value__` 仅适用于单值场景,用于映射度量值,多值场景下,直接指定 field 字段即可
* 必选:是
* 默认值:无
* **metric**
* 描述:仅适用于 RDB 场景下,表示数据迁移任务需要迁移的 Metric 列表
* 必选:否
* 默认值:无
* **field**
* 描述:仅适用于多值场景下,表示数据迁移任务需要迁移的 Field 列表
* 必选:否
* 默认值:无
* **tag**
* 描述:数据迁移任务需要迁移的 TagK 和 TagV,用于进一步过滤时间线
* 必选:否
* 默认值:无
* **splitIntervalMs**
* 描述:用于 DataX 内部切分 Task,每个 Task 只查询一小部分的时间段
* 必选:是
* 默认值:无
* 注意:单位是 ms 毫秒
* **beginDateTime**
* 描述:和 endDateTime 配合使用,用于指定哪个时间段内的数据点,需要被迁移
* 必选:是
* 格式:`yyyy-MM-dd HH:mm:ss`
* 默认值:无
* 注意:指定起止时间会自动忽略分钟和秒,转为整点时刻,例如 2019-4-18 的 [3:35, 4:55) 会被转为 [3:00, 4:00)
* **endDateTime**
* 描述:和 beginDateTime 配合使用,用于指定哪个时间段内的数据点,需要被迁移
* 必选:是
* 格式:`yyyy-MM-dd HH:mm:ss`
* 默认值:无
* 注意:指定起止时间会自动忽略分钟和秒,转为整点时刻,例如 2019-4-18 的 [3:35, 4:55) 会被转为 [3:00, 4:00)
### 3.3 类型转换
| DataX 内部类型 | TSDB 数据类型 |
| -------------- | ------------------------------------------------------------ |
| String | TSDB 数据点序列化字符串,包括 timestamp、metric、tags、fields 和 value |
## 4 约束限制
### 4.2 如果存在某一个 Metric 下在一个小时范围内的数据量过大,可能需要通过 `-j` 参数调整 JVM 内存大小
考虑到下游 Writer 如果写入速度不及 TSDB Reader 的查询数据,可能会存在积压的情况,因此需要适当地调整 JVM 参数。以"从 阿里云 TSDB 数据库同步抽取数据到本地的作业"为例,启动命令如下:
```bash
python datax/bin/datax.py tsdb2stream.json -j "-Xms4096m -Xmx4096m"
```
### 4.3 指定起止时间会自动被转为整点时刻
指定起止时间会自动被转为整点时刻,例如 2019-4-18 的 `[3:35, 3:55)` 会被转为 `[3:00, 4:00)`
+146
View File
@@ -0,0 +1,146 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>com.alibaba.datax</groupId>
<artifactId>datax-all</artifactId>
<version>0.0.1-SNAPSHOT</version>
</parent>
<artifactId>tsdbreader</artifactId>
<name>tsdbreader</name>
<packaging>jar</packaging>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<!-- common -->
<commons-lang3.version>3.3.2</commons-lang3.version>
<!-- http -->
<httpclient.version>4.4</httpclient.version>
<commons-io.version>2.4</commons-io.version>
<!-- json -->
<fastjson.version>1.2.28</fastjson.version>
<!-- test -->
<junit4.version>4.12</junit4.version>
<!-- time -->
<joda-time.version>2.9.9</joda-time.version>
</properties>
<dependencies>
<dependency>
<groupId>com.alibaba.datax</groupId>
<artifactId>datax-common</artifactId>
<version>${datax-project-version}</version>
<exclusions>
<exclusion>
<artifactId>slf4j-log4j12</artifactId>
<groupId>org.slf4j</groupId>
</exclusion>
<exclusion>
<artifactId>fastjson</artifactId>
<groupId>com.alibaba</groupId>
</exclusion>
<exclusion>
<artifactId>commons-math3</artifactId>
<groupId>org.apache.commons</groupId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
</dependency>
<dependency>
<groupId>ch.qos.logback</groupId>
<artifactId>logback-classic</artifactId>
</dependency>
<!-- common -->
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
<version>${commons-lang3.version}</version>
</dependency>
<!-- http -->
<dependency>
<groupId>org.apache.httpcomponents</groupId>
<artifactId>httpclient</artifactId>
<version>${httpclient.version}</version>
</dependency>
<dependency>
<groupId>commons-io</groupId>
<artifactId>commons-io</artifactId>
<version>${commons-io.version}</version>
</dependency>
<dependency>
<groupId>org.apache.httpcomponents</groupId>
<artifactId>fluent-hc</artifactId>
<version>${httpclient.version}</version>
</dependency>
<!-- json -->
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>${fastjson.version}</version>
</dependency>
<!-- time -->
<dependency>
<groupId>joda-time</groupId>
<artifactId>joda-time</artifactId>
<version>${joda-time.version}</version>
</dependency>
<!-- test -->
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<version>${junit4.version}</version>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
<!-- compiler plugin -->
<plugin>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<source>${jdk-version}</source>
<target>${jdk-version}</target>
<encoding>${project-sourceEncoding}</encoding>
</configuration>
</plugin>
<!-- assembly plugin -->
<plugin>
<artifactId>maven-assembly-plugin</artifactId>
<configuration>
<descriptors>
<descriptor>src/main/assembly/package.xml</descriptor>
</descriptors>
<finalName>datax</finalName>
</configuration>
<executions>
<execution>
<id>dwzip</id>
<phase>package</phase>
<goals>
<goal>single</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
+35
View File
@@ -0,0 +1,35 @@
<assembly
xmlns="http://maven.apache.org/plugins/maven-assembly-plugin/assembly/1.1.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/plugins/maven-assembly-plugin/assembly/1.1.0 http://maven.apache.org/xsd/assembly-1.1.0.xsd">
<id></id>
<formats>
<format>dir</format>
</formats>
<includeBaseDirectory>false</includeBaseDirectory>
<fileSets>
<fileSet>
<directory>src/main/resources</directory>
<includes>
<include>plugin.json</include>
<include>plugin_job_template.json</include>
</includes>
<outputDirectory>plugin/reader/tsdbreader</outputDirectory>
</fileSet>
<fileSet>
<directory>target/</directory>
<includes>
<include>tsdbreader-0.0.1-SNAPSHOT.jar</include>
</includes>
<outputDirectory>plugin/reader/tsdbreader</outputDirectory>
</fileSet>
</fileSets>
<dependencySets>
<dependencySet>
<useProjectArtifact>false</useProjectArtifact>
<outputDirectory>plugin/reader/tsdbreader/libs</outputDirectory>
<scope>runtime</scope>
</dependencySet>
</dependencySets>
</assembly>
@@ -0,0 +1,29 @@
package com.alibaba.datax.plugin.reader.tsdbreader;
import java.util.HashSet;
import java.util.Set;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionConstant
*
* @author Benedict Jin
* @since 2019-10-21
*/
public final class Constant {
static final String DEFAULT_DATA_FORMAT = "yyyy-MM-dd HH:mm:ss";
public static final String METRIC_SPECIFY_KEY = "__metric__";
public static final String TS_SPECIFY_KEY = "__ts__";
public static final String VALUE_SPECIFY_KEY = "__value__";
static final Set<String> MUST_CONTAINED_SPECIFY_KEYS = new HashSet<>();
static {
MUST_CONTAINED_SPECIFY_KEYS.add(METRIC_SPECIFY_KEY);
MUST_CONTAINED_SPECIFY_KEYS.add(TS_SPECIFY_KEY);
// __value__ 在多值场景下,可以不指定
}
}
@@ -0,0 +1,36 @@
package com.alibaba.datax.plugin.reader.tsdbreader;
import java.util.HashSet;
import java.util.Set;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionKey
*
* @author Benedict Jin
* @since 2019-10-21
*/
public class Key {
// TSDB for OpenTSDB / InfluxDB / TimeScale / Prometheus etc.
// RDB for MySQL / ADB etc.
static final String SINK_DB_TYPE = "sinkDbType";
static final String ENDPOINT = "endpoint";
static final String COLUMN = "column";
static final String METRIC = "metric";
static final String FIELD = "field";
static final String TAG = "tag";
static final String INTERVAL_DATE_TIME = "splitIntervalMs";
static final String BEGIN_DATE_TIME = "beginDateTime";
static final String END_DATE_TIME = "endDateTime";
static final Integer INTERVAL_DATE_TIME_DEFAULT_VALUE = 60;
static final String TYPE_DEFAULT_VALUE = "TSDB";
static final Set<String> TYPE_SET = new HashSet<>();
static {
TYPE_SET.add("TSDB");
TYPE_SET.add("RDB");
}
}
@@ -0,0 +1,320 @@
package com.alibaba.datax.plugin.reader.tsdbreader;
import com.alibaba.datax.common.exception.DataXException;
import com.alibaba.datax.common.plugin.RecordSender;
import com.alibaba.datax.common.spi.Reader;
import com.alibaba.datax.common.util.Configuration;
import com.alibaba.datax.plugin.reader.tsdbreader.conn.TSDBConnection;
import com.alibaba.datax.plugin.reader.tsdbreader.util.TimeUtils;
import com.alibaba.fastjson.JSON;
import org.apache.commons.lang3.StringUtils;
import org.joda.time.DateTime;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionTSDB Reader
*
* @author Benedict Jin
* @since 2019-10-21
*/
@SuppressWarnings("unused")
public class TSDBReader extends Reader {
public static class Job extends Reader.Job {
private static final Logger LOG = LoggerFactory.getLogger(Job.class);
private Configuration originalConfig;
@Override
public void init() {
this.originalConfig = super.getPluginJobConf();
String type = originalConfig.getString(Key.SINK_DB_TYPE, Key.TYPE_DEFAULT_VALUE);
if (StringUtils.isBlank(type)) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.REQUIRED_VALUE,
"The parameter [" + Key.SINK_DB_TYPE + "] is not set.");
}
if (!Key.TYPE_SET.contains(type)) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.ILLEGAL_VALUE,
"The parameter [" + Key.SINK_DB_TYPE + "] should be one of [" +
JSON.toJSONString(Key.TYPE_SET) + "].");
}
String address = originalConfig.getString(Key.ENDPOINT);
if (StringUtils.isBlank(address)) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.REQUIRED_VALUE,
"The parameter [" + Key.ENDPOINT + "] is not set.");
}
// tagK / field could be empty
if ("TSDB".equals(type)) {
List<String> columns = originalConfig.getList(Key.COLUMN, String.class);
if (columns == null || columns.isEmpty()) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.REQUIRED_VALUE,
"The parameter [" + Key.COLUMN + "] is not set.");
}
} else {
List<String> columns = originalConfig.getList(Key.COLUMN, String.class);
if (columns == null || columns.isEmpty()) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.REQUIRED_VALUE,
"The parameter [" + Key.COLUMN + "] is not set.");
}
for (String specifyKey : Constant.MUST_CONTAINED_SPECIFY_KEYS) {
if (!columns.contains(specifyKey)) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.ILLEGAL_VALUE,
"The parameter [" + Key.COLUMN + "] should contain "
+ JSON.toJSONString(Constant.MUST_CONTAINED_SPECIFY_KEYS) + ".");
}
}
final List<String> metrics = originalConfig.getList(Key.METRIC, String.class);
if (metrics == null || metrics.isEmpty()) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.REQUIRED_VALUE,
"The parameter [" + Key.METRIC + "] is not set.");
}
}
Integer splitIntervalMs = originalConfig.getInt(Key.INTERVAL_DATE_TIME,
Key.INTERVAL_DATE_TIME_DEFAULT_VALUE);
if (splitIntervalMs <= 0) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.ILLEGAL_VALUE,
"The parameter [" + Key.INTERVAL_DATE_TIME + "] should be great than zero.");
}
SimpleDateFormat format = new SimpleDateFormat(Constant.DEFAULT_DATA_FORMAT);
String startTime = originalConfig.getString(Key.BEGIN_DATE_TIME);
Long startDate;
if (startTime == null || startTime.trim().length() == 0) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.REQUIRED_VALUE,
"The parameter [" + Key.BEGIN_DATE_TIME + "] is not set.");
} else {
try {
startDate = format.parse(startTime).getTime();
} catch (ParseException e) {
throw DataXException.asDataXException(TSDBReaderErrorCode.ILLEGAL_VALUE,
"The parameter [" + Key.BEGIN_DATE_TIME +
"] needs to conform to the [" + Constant.DEFAULT_DATA_FORMAT + "] format.");
}
}
String endTime = originalConfig.getString(Key.END_DATE_TIME);
Long endDate;
if (endTime == null || endTime.trim().length() == 0) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.REQUIRED_VALUE,
"The parameter [" + Key.END_DATE_TIME + "] is not set.");
} else {
try {
endDate = format.parse(endTime).getTime();
} catch (ParseException e) {
throw DataXException.asDataXException(TSDBReaderErrorCode.ILLEGAL_VALUE,
"The parameter [" + Key.END_DATE_TIME +
"] needs to conform to the [" + Constant.DEFAULT_DATA_FORMAT + "] format.");
}
}
if (startDate >= endDate) {
throw DataXException.asDataXException(TSDBReaderErrorCode.ILLEGAL_VALUE,
"The parameter [" + Key.BEGIN_DATE_TIME +
"] should be less than the parameter [" + Key.END_DATE_TIME + "].");
}
}
@Override
public void prepare() {
}
@Override
public List<Configuration> split(int adviceNumber) {
List<Configuration> configurations = new ArrayList<>();
// get metrics
String type = originalConfig.getString(Key.SINK_DB_TYPE, Key.TYPE_DEFAULT_VALUE);
List<String> columns4TSDB = null;
List<String> columns4RDB = null;
List<String> metrics = null;
if ("TSDB".equals(type)) {
columns4TSDB = originalConfig.getList(Key.COLUMN, String.class);
} else {
columns4RDB = originalConfig.getList(Key.COLUMN, String.class);
metrics = originalConfig.getList(Key.METRIC, String.class);
}
// get time interval
Integer splitIntervalMs = originalConfig.getInt(Key.INTERVAL_DATE_TIME,
Key.INTERVAL_DATE_TIME_DEFAULT_VALUE);
// get time range
SimpleDateFormat format = new SimpleDateFormat(Constant.DEFAULT_DATA_FORMAT);
long startTime;
try {
startTime = format.parse(originalConfig.getString(Key.BEGIN_DATE_TIME)).getTime();
} catch (ParseException e) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.ILLEGAL_VALUE, "解析[" + Key.BEGIN_DATE_TIME + "]失败.", e);
}
long endTime;
try {
endTime = format.parse(originalConfig.getString(Key.END_DATE_TIME)).getTime();
} catch (ParseException e) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.ILLEGAL_VALUE, "解析[" + Key.END_DATE_TIME + "]失败.", e);
}
if (TimeUtils.isSecond(startTime)) {
startTime *= 1000;
}
if (TimeUtils.isSecond(endTime)) {
endTime *= 1000;
}
DateTime startDateTime = new DateTime(TimeUtils.getTimeInHour(startTime));
DateTime endDateTime = new DateTime(TimeUtils.getTimeInHour(endTime));
if ("TSDB".equals(type)) {
// split by metric
for (String column : columns4TSDB) {
// split by time in hour
while (startDateTime.isBefore(endDateTime)) {
Configuration clone = this.originalConfig.clone();
clone.set(Key.COLUMN, Collections.singletonList(column));
clone.set(Key.BEGIN_DATE_TIME, startDateTime.getMillis());
startDateTime = startDateTime.plusMillis(splitIntervalMs);
// Make sure the time interval is [start, end).
clone.set(Key.END_DATE_TIME, startDateTime.getMillis() - 1);
configurations.add(clone);
LOG.info("Configuration: {}", JSON.toJSONString(clone));
}
}
} else {
// split by metric
for (String metric : metrics) {
// split by time in hour
while (startDateTime.isBefore(endDateTime)) {
Configuration clone = this.originalConfig.clone();
clone.set(Key.COLUMN, columns4RDB);
clone.set(Key.METRIC, Collections.singletonList(metric));
clone.set(Key.BEGIN_DATE_TIME, startDateTime.getMillis());
startDateTime = startDateTime.plusMillis(splitIntervalMs);
// Make sure the time interval is [start, end).
clone.set(Key.END_DATE_TIME, startDateTime.getMillis() - 1);
configurations.add(clone);
LOG.info("Configuration: {}", JSON.toJSONString(clone));
}
}
}
return configurations;
}
@Override
public void post() {
}
@Override
public void destroy() {
}
}
public static class Task extends Reader.Task {
private static final Logger LOG = LoggerFactory.getLogger(Task.class);
private String type;
private List<String> columns4TSDB = null;
private List<String> columns4RDB = null;
private List<String> metrics = null;
private Map<String, Object> fields;
private Map<String, Object> tags;
private TSDBConnection conn;
private Long startTime;
private Long endTime;
@Override
public void init() {
Configuration readerSliceConfig = super.getPluginJobConf();
LOG.info("getPluginJobConf: {}", JSON.toJSONString(readerSliceConfig));
this.type = readerSliceConfig.getString(Key.SINK_DB_TYPE);
if ("TSDB".equals(type)) {
columns4TSDB = readerSliceConfig.getList(Key.COLUMN, String.class);
} else {
columns4RDB = readerSliceConfig.getList(Key.COLUMN, String.class);
metrics = readerSliceConfig.getList(Key.METRIC, String.class);
}
this.fields = readerSliceConfig.getMap(Key.FIELD);
this.tags = readerSliceConfig.getMap(Key.TAG);
String address = readerSliceConfig.getString(Key.ENDPOINT);
conn = new TSDBConnection(address);
this.startTime = readerSliceConfig.getLong(Key.BEGIN_DATE_TIME);
this.endTime = readerSliceConfig.getLong(Key.END_DATE_TIME);
}
@Override
public void prepare() {
}
@Override
@SuppressWarnings("unchecked")
public void startRead(RecordSender recordSender) {
try {
if ("TSDB".equals(type)) {
for (String metric : columns4TSDB) {
final Map<String, String> tags = this.tags == null ?
null : (Map<String, String>) this.tags.get(metric);
if (fields == null || !fields.containsKey(metric)) {
conn.sendDPs(metric, tags, this.startTime, this.endTime, recordSender);
} else {
conn.sendDPs(metric, (List<String>) fields.get(metric),
tags, this.startTime, this.endTime, recordSender);
}
}
} else {
for (String metric : metrics) {
final Map<String, String> tags = this.tags == null ?
null : (Map<String, String>) this.tags.get(metric);
if (fields == null || !fields.containsKey(metric)) {
conn.sendRecords(metric, tags, startTime, endTime, columns4RDB, recordSender);
} else {
conn.sendRecords(metric, (List<String>) fields.get(metric),
tags, startTime, endTime, columns4RDB, recordSender);
}
}
}
} catch (Exception e) {
throw DataXException.asDataXException(
TSDBReaderErrorCode.ILLEGAL_VALUE, "获取或发送数据点的过程中出错!", e);
}
}
@Override
public void post() {
}
@Override
public void destroy() {
}
}
}
@@ -0,0 +1,40 @@
package com.alibaba.datax.plugin.reader.tsdbreader;
import com.alibaba.datax.common.spi.ErrorCode;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionTSDB Reader Error Code
*
* @author Benedict Jin
* @since 2019-10-21
*/
public enum TSDBReaderErrorCode implements ErrorCode {
REQUIRED_VALUE("TSDBReader-00", "缺失必要的值"),
ILLEGAL_VALUE("TSDBReader-01", "值非法");
private final String code;
private final String description;
TSDBReaderErrorCode(String code, String description) {
this.code = code;
this.description = description;
}
@Override
public String getCode() {
return this.code;
}
@Override
public String getDescription() {
return this.description;
}
@Override
public String toString() {
return String.format("Code:[%s], Description:[%s]. ", this.code, this.description);
}
}
@@ -0,0 +1,88 @@
package com.alibaba.datax.plugin.reader.tsdbreader.conn;
import com.alibaba.datax.common.plugin.RecordSender;
import java.util.List;
import java.util.Map;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionConnection for TSDB-like databases
*
* @author Benedict Jin
* @since 2019-10-21
*/
public interface Connection4TSDB {
/**
* Get the address of Database.
*
* @return host+ip
*/
String address();
/**
* Get the version of Database.
*
* @return version
*/
String version();
/**
* Get these configurations.
*
* @return configs
*/
String config();
/**
* Get the list of supported version.
*
* @return version list
*/
String[] getSupportVersionPrefix();
/**
* Send data points for TSDB with single field.
*/
void sendDPs(String metric, Map<String, String> tags, Long start, Long end, RecordSender recordSender) throws Exception;
/**
* Send data points for TSDB with multi fields.
*/
void sendDPs(String metric, List<String> fields, Map<String, String> tags, Long start, Long end, RecordSender recordSender) throws Exception;
/**
* Send data points for RDB with single field.
*/
void sendRecords(String metric, Map<String, String> tags, Long start, Long end, List<String> columns4RDB, RecordSender recordSender) throws Exception;
/**
* Send data points for RDB with multi fields.
*/
void sendRecords(String metric, List<String> fields, Map<String, String> tags, Long start, Long end, List<String> columns4RDB, RecordSender recordSender) throws Exception;
/**
* Put data point.
*
* @param dp data point
* @return whether the data point is written successfully
*/
boolean put(DataPoint4TSDB dp);
/**
* Put data points.
*
* @param dps data points
* @return whether the data point is written successfully
*/
boolean put(List<DataPoint4TSDB> dps);
/**
* Whether current version is supported.
*
* @return true: supported; false: not yet!
*/
boolean isSupported();
}
@@ -0,0 +1,68 @@
package com.alibaba.datax.plugin.reader.tsdbreader.conn;
import com.alibaba.fastjson.JSON;
import java.util.Map;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionDataPoint for TSDB with Multi Fields
*
* @author Benedict Jin
* @since 2019-10-21
*/
public class DataPoint4MultiFieldsTSDB {
private long timestamp;
private String metric;
private Map<String, Object> tags;
private Map<String, Object> fields;
public DataPoint4MultiFieldsTSDB() {
}
public DataPoint4MultiFieldsTSDB(long timestamp, String metric, Map<String, Object> tags, Map<String, Object> fields) {
this.timestamp = timestamp;
this.metric = metric;
this.tags = tags;
this.fields = fields;
}
public long getTimestamp() {
return timestamp;
}
public void setTimestamp(long timestamp) {
this.timestamp = timestamp;
}
public String getMetric() {
return metric;
}
public void setMetric(String metric) {
this.metric = metric;
}
public Map<String, Object> getTags() {
return tags;
}
public void setTags(Map<String, Object> tags) {
this.tags = tags;
}
public Map<String, Object> getFields() {
return fields;
}
public void setFields(Map<String, Object> fields) {
this.fields = fields;
}
@Override
public String toString() {
return JSON.toJSONString(this);
}
}
@@ -0,0 +1,68 @@
package com.alibaba.datax.plugin.reader.tsdbreader.conn;
import com.alibaba.fastjson.JSON;
import java.util.Map;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionDataPoint for TSDB
*
* @author Benedict Jin
* @since 2019-10-21
*/
public class DataPoint4TSDB {
private long timestamp;
private String metric;
private Map<String, Object> tags;
private Object value;
public DataPoint4TSDB() {
}
public DataPoint4TSDB(long timestamp, String metric, Map<String, Object> tags, Object value) {
this.timestamp = timestamp;
this.metric = metric;
this.tags = tags;
this.value = value;
}
public long getTimestamp() {
return timestamp;
}
public void setTimestamp(long timestamp) {
this.timestamp = timestamp;
}
public String getMetric() {
return metric;
}
public void setMetric(String metric) {
this.metric = metric;
}
public Map<String, Object> getTags() {
return tags;
}
public void setTags(Map<String, Object> tags) {
this.tags = tags;
}
public Object getValue() {
return value;
}
public void setValue(Object value) {
this.value = value;
}
@Override
public String toString() {
return JSON.toJSONString(this);
}
}
@@ -0,0 +1,64 @@
package com.alibaba.datax.plugin.reader.tsdbreader.conn;
import java.util.List;
import java.util.Map;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionMulti Field Query Result
*
* @author Benedict Jin
* @since 2019-10-22
*/
public class MultiFieldQueryResult {
private String metric;
private Map<String, Object> tags;
private List<String> aggregatedTags;
private List<String> columns;
private List<List<Object>> values;
public MultiFieldQueryResult() {
}
public String getMetric() {
return metric;
}
public void setMetric(String metric) {
this.metric = metric;
}
public Map<String, Object> getTags() {
return tags;
}
public void setTags(Map<String, Object> tags) {
this.tags = tags;
}
public List<String> getAggregatedTags() {
return aggregatedTags;
}
public void setAggregatedTags(List<String> aggregatedTags) {
this.aggregatedTags = aggregatedTags;
}
public List<String> getColumns() {
return columns;
}
public void setColumns(List<String> columns) {
this.columns = columns;
}
public List<List<Object>> getValues() {
return values;
}
public void setValues(List<List<Object>> values) {
this.values = values;
}
}
@@ -0,0 +1,64 @@
package com.alibaba.datax.plugin.reader.tsdbreader.conn;
import java.util.List;
import java.util.Map;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionQuery Result
*
* @author Benedict Jin
* @since 2019-09-19
*/
public class QueryResult {
private String metricName;
private Map<String, Object> tags;
private List<String> groupByTags;
private List<String> aggregatedTags;
private Map<String, Object> dps;
public QueryResult() {
}
public String getMetricName() {
return metricName;
}
public void setMetricName(String metricName) {
this.metricName = metricName;
}
public Map<String, Object> getTags() {
return tags;
}
public void setTags(Map<String, Object> tags) {
this.tags = tags;
}
public List<String> getGroupByTags() {
return groupByTags;
}
public void setGroupByTags(List<String> groupByTags) {
this.groupByTags = groupByTags;
}
public List<String> getAggregatedTags() {
return aggregatedTags;
}
public void setAggregatedTags(List<String> aggregatedTags) {
this.aggregatedTags = aggregatedTags;
}
public Map<String, Object> getDps() {
return dps;
}
public void setDps(Map<String, Object> dps) {
this.dps = dps;
}
}
@@ -0,0 +1,94 @@
package com.alibaba.datax.plugin.reader.tsdbreader.conn;
import com.alibaba.datax.common.plugin.RecordSender;
import com.alibaba.datax.plugin.reader.tsdbreader.util.TSDBUtils;
import com.alibaba.fastjson.JSON;
import org.apache.commons.lang3.StringUtils;
import java.util.List;
import java.util.Map;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionTSDB Connection
*
* @author Benedict Jin
* @since 2019-10-21
*/
public class TSDBConnection implements Connection4TSDB {
private String address;
public TSDBConnection(String address) {
this.address = address;
}
@Override
public String address() {
return address;
}
@Override
public String version() {
return TSDBUtils.version(address);
}
@Override
public String config() {
return TSDBUtils.config(address);
}
@Override
public String[] getSupportVersionPrefix() {
return new String[]{"2.4", "2.5"};
}
@Override
public void sendDPs(String metric, Map<String, String> tags, Long start, Long end, RecordSender recordSender) throws Exception {
TSDBDump.dump4TSDB(this, metric, tags, start, end, recordSender);
}
@Override
public void sendDPs(String metric, List<String> fields, Map<String, String> tags, Long start, Long end, RecordSender recordSender) throws Exception {
TSDBDump.dump4TSDB(this, metric, fields, tags, start, end, recordSender);
}
@Override
public void sendRecords(String metric, Map<String, String> tags, Long start, Long end, List<String> columns4RDB, RecordSender recordSender) throws Exception {
TSDBDump.dump4RDB(this, metric, tags, start, end, columns4RDB, recordSender);
}
@Override
public void sendRecords(String metric, List<String> fields, Map<String, String> tags, Long start, Long end, List<String> columns4RDB, RecordSender recordSender) throws Exception {
TSDBDump.dump4RDB(this, metric, fields, tags, start, end, columns4RDB, recordSender);
}
@Override
public boolean put(DataPoint4TSDB dp) {
return false;
}
@Override
public boolean put(List<DataPoint4TSDB> dps) {
return false;
}
@Override
public boolean isSupported() {
String versionJson = version();
if (StringUtils.isBlank(versionJson)) {
throw new RuntimeException("Cannot get the version!");
}
String version = JSON.parseObject(versionJson).getString("version");
if (StringUtils.isBlank(version)) {
return false;
}
for (String prefix : getSupportVersionPrefix()) {
if (version.startsWith(prefix)) {
return true;
}
}
return false;
}
}
@@ -0,0 +1,318 @@
package com.alibaba.datax.plugin.reader.tsdbreader.conn;
import com.alibaba.datax.common.element.*;
import com.alibaba.datax.common.plugin.RecordSender;
import com.alibaba.datax.plugin.reader.tsdbreader.Constant;
import com.alibaba.datax.plugin.reader.tsdbreader.util.HttpUtils;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.parser.Feature;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.HashMap;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionTSDB Dump
*
* @author Benedict Jin
* @since 2019-10-21
*/
final class TSDBDump {
private static final Logger LOG = LoggerFactory.getLogger(TSDBDump.class);
private static final String QUERY = "/api/query";
private static final String QUERY_MULTI_FIELD = "/api/mquery";
static {
JSON.DEFAULT_PARSER_FEATURE &= ~Feature.UseBigDecimal.getMask();
}
private TSDBDump() {
}
static void dump4TSDB(TSDBConnection conn, String metric, Map<String, String> tags,
Long start, Long end, RecordSender sender) throws Exception {
LOG.info("conn address: {}, metric: {}, start: {}, end: {}", conn.address(), metric, start, end);
String res = queryRange4SingleField(conn, metric, tags, start, end);
List<String> dps = getDps4TSDB(metric, res);
if (dps == null || dps.isEmpty()) {
return;
}
sendTSDBDps(sender, dps);
}
static void dump4TSDB(TSDBConnection conn, String metric, List<String> fields, Map<String, String> tags,
Long start, Long end, RecordSender sender) throws Exception {
LOG.info("conn address: {}, metric: {}, start: {}, end: {}", conn.address(), metric, start, end);
String res = queryRange4MultiFields(conn, metric, fields, tags, start, end);
List<String> dps = getDps4TSDB(metric, fields, res);
if (dps == null || dps.isEmpty()) {
return;
}
sendTSDBDps(sender, dps);
}
static void dump4RDB(TSDBConnection conn, String metric, Map<String, String> tags,
Long start, Long end, List<String> columns4RDB, RecordSender sender) throws Exception {
LOG.info("conn address: {}, metric: {}, start: {}, end: {}", conn.address(), metric, start, end);
String res = queryRange4SingleField(conn, metric, tags, start, end);
List<DataPoint4TSDB> dps = getDps4RDB(metric, res);
if (dps == null || dps.isEmpty()) {
return;
}
for (DataPoint4TSDB dp : dps) {
final Record record = sender.createRecord();
final Map<String, Object> tagKV = dp.getTags();
for (String column : columns4RDB) {
if (Constant.METRIC_SPECIFY_KEY.equals(column)) {
record.addColumn(new StringColumn(dp.getMetric()));
} else if (Constant.TS_SPECIFY_KEY.equals(column)) {
record.addColumn(new LongColumn(dp.getTimestamp()));
} else if (Constant.VALUE_SPECIFY_KEY.equals(column)) {
record.addColumn(getColumn(dp.getValue()));
} else {
final Object tagk = tagKV.get(column);
if (tagk == null) {
continue;
}
record.addColumn(getColumn(tagk));
}
}
sender.sendToWriter(record);
}
}
static void dump4RDB(TSDBConnection conn, String metric, List<String> fields,
Map<String, String> tags, Long start, Long end,
List<String> columns4RDB, RecordSender sender) throws Exception {
LOG.info("conn address: {}, metric: {}, start: {}, end: {}", conn.address(), metric, start, end);
String res = queryRange4MultiFields(conn, metric, fields, tags, start, end);
List<DataPoint4TSDB> dps = getDps4RDB(metric, fields, res);
if (dps == null || dps.isEmpty()) {
return;
}
for (DataPoint4TSDB dp : dps) {
final Record record = sender.createRecord();
final Map<String, Object> tagKV = dp.getTags();
for (String column : columns4RDB) {
if (Constant.METRIC_SPECIFY_KEY.equals(column)) {
record.addColumn(new StringColumn(dp.getMetric()));
} else if (Constant.TS_SPECIFY_KEY.equals(column)) {
record.addColumn(new LongColumn(dp.getTimestamp()));
} else {
final Object tagvOrField = tagKV.get(column);
if (tagvOrField == null) {
continue;
}
record.addColumn(getColumn(tagvOrField));
}
}
sender.sendToWriter(record);
}
}
private static Column getColumn(Object value) throws Exception {
Column valueColumn;
if (value instanceof Double) {
valueColumn = new DoubleColumn((Double) value);
} else if (value instanceof Long) {
valueColumn = new LongColumn((Long) value);
} else if (value instanceof String) {
valueColumn = new StringColumn((String) value);
} else {
throw new Exception(String.format("value 不支持类型: [%s]", value.getClass().getSimpleName()));
}
return valueColumn;
}
private static String queryRange4SingleField(TSDBConnection conn, String metric, Map<String, String> tags,
Long start, Long end) throws Exception {
String tagKV = getFilterByTags(tags);
String body = "{\n" +
" \"start\": " + start + ",\n" +
" \"end\": " + end + ",\n" +
" \"queries\": [\n" +
" {\n" +
" \"aggregator\": \"none\",\n" +
" \"metric\": \"" + metric + "\"\n" +
(tagKV == null ? "" : tagKV) +
" }\n" +
" ]\n" +
"}";
return HttpUtils.post(conn.address() + QUERY, body);
}
private static String queryRange4MultiFields(TSDBConnection conn, String metric, List<String> fields,
Map<String, String> tags, Long start, Long end) throws Exception {
// fields
StringBuilder fieldBuilder = new StringBuilder();
fieldBuilder.append("\"fields\":[");
for (int i = 0; i < fields.size(); i++) {
fieldBuilder.append("{\"field\": \"").append(fields.get(i)).append("\",\"aggregator\": \"none\"}");
if (i != fields.size() - 1) {
fieldBuilder.append(",");
}
}
fieldBuilder.append("]");
// tagkv
String tagKV = getFilterByTags(tags);
String body = "{\n" +
" \"start\": " + start + ",\n" +
" \"end\": " + end + ",\n" +
" \"queries\": [\n" +
" {\n" +
" \"aggregator\": \"none\",\n" +
" \"metric\": \"" + metric + "\",\n" +
fieldBuilder.toString() +
(tagKV == null ? "" : tagKV) +
" }\n" +
" ]\n" +
"}";
return HttpUtils.post(conn.address() + QUERY_MULTI_FIELD, body);
}
private static String getFilterByTags(Map<String, String> tags) {
if (tags != null && !tags.isEmpty()) {
// tagKV = ",\"tags:\":" + JSON.toJSONString(tags);
StringBuilder tagBuilder = new StringBuilder();
tagBuilder.append(",\"filters\":[");
int count = 1;
final int size = tags.size();
for (Map.Entry<String, String> entry : tags.entrySet()) {
final String tagK = entry.getKey();
final String tagV = entry.getValue();
tagBuilder.append("{\"type\":\"literal_or\",\"tagk\":\"").append(tagK)
.append("\",\"filter\":\"").append(tagV).append("\",\"groupBy\":false}");
if (count != size) {
tagBuilder.append(",");
}
count++;
}
tagBuilder.append("]");
return tagBuilder.toString();
}
return null;
}
private static List<String> getDps4TSDB(String metric, String dps) {
final List<QueryResult> jsonArray = JSON.parseArray(dps, QueryResult.class);
if (jsonArray.size() == 0) {
return null;
}
List<String> dpsArr = new LinkedList<>();
for (QueryResult queryResult : jsonArray) {
final Map<String, Object> tags = queryResult.getTags();
final Map<String, Object> points = queryResult.getDps();
for (Map.Entry<String, Object> entry : points.entrySet()) {
final String ts = entry.getKey();
final Object value = entry.getValue();
DataPoint4TSDB dp = new DataPoint4TSDB();
dp.setMetric(metric);
dp.setTags(tags);
dp.setTimestamp(Long.parseLong(ts));
dp.setValue(value);
dpsArr.add(dp.toString());
}
}
return dpsArr;
}
private static List<String> getDps4TSDB(String metric, List<String> fields, String dps) {
final List<MultiFieldQueryResult> jsonArray = JSON.parseArray(dps, MultiFieldQueryResult.class);
if (jsonArray.size() == 0) {
return null;
}
List<String> dpsArr = new LinkedList<>();
for (MultiFieldQueryResult queryResult : jsonArray) {
final Map<String, Object> tags = queryResult.getTags();
final List<List<Object>> values = queryResult.getValues();
for (List<Object> value : values) {
final String ts = value.get(0).toString();
Map<String, Object> fieldsAndValues = new HashMap<>();
for (int i = 0; i < fields.size(); i++) {
fieldsAndValues.put(fields.get(i), value.get(i + 1));
}
final DataPoint4MultiFieldsTSDB dp = new DataPoint4MultiFieldsTSDB();
dp.setMetric(metric);
dp.setTimestamp(Long.parseLong(ts));
dp.setTags(tags);
dp.setFields(fieldsAndValues);
dpsArr.add(dp.toString());
}
}
return dpsArr;
}
private static List<DataPoint4TSDB> getDps4RDB(String metric, String dps) {
final List<QueryResult> jsonArray = JSON.parseArray(dps, QueryResult.class);
if (jsonArray.size() == 0) {
return null;
}
List<DataPoint4TSDB> dpsArr = new LinkedList<>();
for (QueryResult queryResult : jsonArray) {
final Map<String, Object> tags = queryResult.getTags();
final Map<String, Object> points = queryResult.getDps();
for (Map.Entry<String, Object> entry : points.entrySet()) {
final String ts = entry.getKey();
final Object value = entry.getValue();
final DataPoint4TSDB dp = new DataPoint4TSDB();
dp.setMetric(metric);
dp.setTags(tags);
dp.setTimestamp(Long.parseLong(ts));
dp.setValue(value);
dpsArr.add(dp);
}
}
return dpsArr;
}
private static List<DataPoint4TSDB> getDps4RDB(String metric, List<String> fields, String dps) {
final List<MultiFieldQueryResult> jsonArray = JSON.parseArray(dps, MultiFieldQueryResult.class);
if (jsonArray.size() == 0) {
return null;
}
List<DataPoint4TSDB> dpsArr = new LinkedList<>();
for (MultiFieldQueryResult queryResult : jsonArray) {
final Map<String, Object> tags = queryResult.getTags();
final List<List<Object>> values = queryResult.getValues();
for (List<Object> value : values) {
final String ts = value.get(0).toString();
Map<String, Object> tagsTmp = new HashMap<>(tags);
for (int i = 0; i < fields.size(); i++) {
tagsTmp.put(fields.get(i), value.get(i + 1));
}
final DataPoint4TSDB dp = new DataPoint4TSDB();
dp.setMetric(metric);
dp.setTimestamp(Long.parseLong(ts));
dp.setTags(tagsTmp);
dpsArr.add(dp);
}
}
return dpsArr;
}
private static void sendTSDBDps(RecordSender sender, List<String> dps) {
for (String dp : dps) {
StringColumn tsdbColumn = new StringColumn(dp);
Record record = sender.createRecord();
record.addColumn(tsdbColumn);
sender.sendToWriter(record);
}
}
}
@@ -0,0 +1,67 @@
package com.alibaba.datax.plugin.reader.tsdbreader.util;
import com.alibaba.fastjson.JSON;
import org.apache.http.client.fluent.Content;
import org.apache.http.client.fluent.Request;
import org.apache.http.entity.ContentType;
import java.nio.charset.StandardCharsets;
import java.util.Map;
import java.util.concurrent.TimeUnit;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionHttpUtils
*
* @author Benedict Jin
* @since 2019-10-21
*/
public final class HttpUtils {
public final static int CONNECT_TIMEOUT_DEFAULT_IN_MILL = (int) TimeUnit.SECONDS.toMillis(60);
public final static int SOCKET_TIMEOUT_DEFAULT_IN_MILL = (int) TimeUnit.SECONDS.toMillis(60);
private HttpUtils() {
}
public static String get(String url) throws Exception {
Content content = Request.Get(url)
.connectTimeout(CONNECT_TIMEOUT_DEFAULT_IN_MILL)
.socketTimeout(SOCKET_TIMEOUT_DEFAULT_IN_MILL)
.execute()
.returnContent();
if (content == null) {
return null;
}
return content.asString(StandardCharsets.UTF_8);
}
public static String post(String url, Map<String, Object> params) throws Exception {
return post(url, JSON.toJSONString(params), CONNECT_TIMEOUT_DEFAULT_IN_MILL, SOCKET_TIMEOUT_DEFAULT_IN_MILL);
}
public static String post(String url, String params) throws Exception {
return post(url, params, CONNECT_TIMEOUT_DEFAULT_IN_MILL, SOCKET_TIMEOUT_DEFAULT_IN_MILL);
}
public static String post(String url, Map<String, Object> params,
int connectTimeoutInMill, int socketTimeoutInMill) throws Exception {
return post(url, JSON.toJSONString(params), connectTimeoutInMill, socketTimeoutInMill);
}
public static String post(String url, String params,
int connectTimeoutInMill, int socketTimeoutInMill) throws Exception {
Content content = Request.Post(url)
.connectTimeout(connectTimeoutInMill)
.socketTimeout(socketTimeoutInMill)
.addHeader("Content-Type", "application/json")
.bodyString(params, ContentType.APPLICATION_JSON)
.execute()
.returnContent();
if (content == null) {
return null;
}
return content.asString(StandardCharsets.UTF_8);
}
}
@@ -0,0 +1,68 @@
package com.alibaba.datax.plugin.reader.tsdbreader.util;
import com.alibaba.datax.plugin.reader.tsdbreader.conn.DataPoint4TSDB;
import com.alibaba.fastjson.JSON;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionTSDB Utils
*
* @author Benedict Jin
* @since 2019-10-21
*/
public final class TSDBUtils {
private static final Logger LOGGER = LoggerFactory.getLogger(TSDBUtils.class);
private TSDBUtils() {
}
public static String version(String address) {
String url = String.format("%s/api/version", address);
String rsp;
try {
rsp = HttpUtils.get(url);
} catch (Exception e) {
throw new RuntimeException(e);
}
return rsp;
}
public static String config(String address) {
String url = String.format("%s/api/config", address);
String rsp;
try {
rsp = HttpUtils.get(url);
} catch (Exception e) {
throw new RuntimeException(e);
}
return rsp;
}
public static boolean put(String address, List<DataPoint4TSDB> dps) {
return put(address, JSON.toJSON(dps));
}
public static boolean put(String address, DataPoint4TSDB dp) {
return put(address, JSON.toJSON(dp));
}
private static boolean put(String address, Object o) {
String url = String.format("%s/api/put", address);
String rsp;
try {
rsp = HttpUtils.post(url, o.toString());
// If successful, the returned content should be null.
assert rsp == null;
} catch (Exception e) {
LOGGER.error("Address: {}, DataPoints: {}", url, o);
throw new RuntimeException(e);
}
return true;
}
}
@@ -0,0 +1,38 @@
package com.alibaba.datax.plugin.reader.tsdbreader.util;
import java.util.concurrent.TimeUnit;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionTimeUtils
*
* @author Benedict Jin
* @since 2019-10-21
*/
public final class TimeUtils {
private TimeUtils() {
}
private static final long SECOND_MASK = 0xFFFFFFFF00000000L;
private static final long HOUR_IN_MILL = TimeUnit.HOURS.toMillis(1);
/**
* Weather the timestamp is second.
*
* @param ts timestamp
*/
public static boolean isSecond(long ts) {
return (ts & SECOND_MASK) == 0;
}
/**
* Get the hour.
*
* @param ms time in millisecond
*/
public static long getTimeInHour(long ms) {
return ms - ms % HOUR_IN_MILL;
}
}
+10
View File
@@ -0,0 +1,10 @@
{
"name": "tsdbreader",
"class": "com.alibaba.datax.plugin.reader.tsdbreader.TSDBReader",
"description": {
"useScene": "从 TSDB 中摄取数据点",
"mechanism": "通过 /api/query 接口查询出符合条件的数据点",
"warn": "指定起止时间会自动忽略分钟和秒,转为整点时刻,例如 2019-4-18 的 [3:35, 4:55) 会被转为 [3:00, 4:00)"
},
"developer": "Benedict Jin"
}
@@ -0,0 +1,29 @@
{
"name": "tsdbreader",
"parameter": {
"sinkDbType": "RDB",
"endpoint": "http://localhost:8242",
"column": [
"__metric__",
"__ts__",
"app",
"cluster",
"group",
"ip",
"zone",
"__value__"
],
"metric": [
"m"
],
"tag": {
"m": {
"app": "a1",
"cluster": "c1"
}
},
"splitIntervalMs": 60000,
"beginDateTime": "2019-01-01 00:00:00",
"endDateTime": "2019-01-01 01:00:00"
}
}
@@ -0,0 +1,30 @@
package com.alibaba.datax.plugin.reader.tsdbreader.conn;
import org.junit.Assert;
import org.junit.Ignore;
import org.junit.Test;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionTSDB Connection4TSDB Test
*
* @author Benedict Jin
* @since 2019-10-21
*/
@Ignore
public class TSDBConnectionTest {
private static final String TSDB_ADDRESS = "http://localhost:8242";
@Test
public void testVersion() {
String version = new TSDBConnection(TSDB_ADDRESS).version();
Assert.assertNotNull(version);
}
@Test
public void testIsSupported() {
Assert.assertTrue(new TSDBConnection(TSDB_ADDRESS).isSupported());
}
}
@@ -0,0 +1,17 @@
package com.alibaba.datax.plugin.reader.tsdbreader.util;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionConst
*
* @author Benedict Jin
* @since 2019-10-21
*/
final class Const {
private Const() {
}
static final String TSDB_ADDRESS = "http://localhost:8242";
}
@@ -0,0 +1,39 @@
package com.alibaba.datax.plugin.reader.tsdbreader.util;
import org.junit.Assert;
import org.junit.Ignore;
import org.junit.Test;
import java.util.HashMap;
import java.util.Map;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* FunctionHttpUtils Test
*
* @author Benedict Jin
* @since 2019-10-21
*/
@Ignore
public class HttpUtilsTest {
@Test
public void testSimpleCase() throws Exception {
String url = "https://httpbin.org/post";
Map<String, Object> params = new HashMap<>();
params.put("foo", "bar");
String rsp = HttpUtils.post(url, params);
System.out.println(rsp);
Assert.assertNotNull(rsp);
}
@Test
public void testGet() throws Exception {
String url = String.format("%s/api/version", Const.TSDB_ADDRESS);
String rsp = HttpUtils.get(url);
System.out.println(rsp);
Assert.assertNotNull(rsp);
}
}
@@ -0,0 +1,33 @@
package com.alibaba.datax.plugin.reader.tsdbreader.util;
import org.junit.Assert;
import org.junit.Test;
import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.util.Date;
/**
* Copyright @ 2019 alibaba.com
* All right reserved.
* Functioncom.alibaba.datax.common.util
*
* @author Benedict Jin
* @since 2019-10-21
*/
public class TimeUtilsTest {
@Test
public void testIsSecond() {
Assert.assertFalse(TimeUtils.isSecond(System.currentTimeMillis()));
Assert.assertTrue(TimeUtils.isSecond(System.currentTimeMillis() / 1000));
}
@Test
public void testGetTimeInHour() throws ParseException {
SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
Date date = sdf.parse("2019-04-18 15:32:33");
long timeInHour = TimeUtils.getTimeInHour(date.getTime());
Assert.assertEquals("2019-04-18 15:00:00", sdf.format(timeInHour));
}
}
-1
View File
@@ -187,7 +187,6 @@ TxtFileWriter实现了从DataX协议转为本地TXT文件功能,本地文件
| DataX 内部类型| 本地文件 数据类型 |
| -------- | ----- |
|
| Long |Long |
| Double |Double|
| String |String|
+1 -1
View File
@@ -64,7 +64,7 @@ DataX本身作为数据同步框架,将不同数据源的同步抽象为从源
* 配置示例:从stream读取数据并打印到控制台
* 第一步、创建业的配置文件(json格式)
* 第一步、创建业的配置文件(json格式)
可以通过命令查看配置模板: python datax.py -r {YOUR_READER} -w {YOUR_WRITER}