Skip to content
项目
群组
代码片段
帮助
当前项目
正在载入...
登录 / 注册
切换导航面板
S
seatunnel-web
项目
项目
详情
活动
周期分析
仓库
仓库
文件
提交
分支
标签
贡献者
图表
比较
统计图
议题
0
议题
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
CI / CD
CI / CD
流水线
作业
日程
统计图
Wiki
Wiki
代码片段
代码片段
成员
成员
折叠边栏
关闭边栏
活动
图像
聊天
创建新问题
作业
提交
问题看板
Open sidebar
宋勇
seatunnel-web
Commits
d18b47a1
提交
d18b47a1
authored
4月 13, 2024
作者:
宋勇
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
atasource-jdbc-demeng
datasource-jdbc-access datasource-http datasource-xml datasource-csv datasource-excel 增加
上级
373502c0
全部展开
隐藏空白字符变更
内嵌
并排
正在显示
30 个修改的文件
包含
125 行增加
和
112 行删除
+125
-112
pom.xml
...ource/seatunnel-datasource-plugins/datasource-csv/pom.xml
+1
-2
CSVAConfiguration.java
...he/seatunnel/datasource/plugin/csv/CSVAConfiguration.java
+5
-3
CSVDataSourceFactory.java
...seatunnel/datasource/plugin/csv/CSVDataSourceFactory.java
+4
-3
CSVDatasourceChannel.java
...seatunnel/datasource/plugin/csv/CSVDatasourceChannel.java
+10
-8
CSVOptionRule.java
...apache/seatunnel/datasource/plugin/csv/CSVOptionRule.java
+0
-2
pom.xml
...rce/seatunnel-datasource-plugins/datasource-excel/pom.xml
+1
-2
ExcelAConfiguration.java
...eatunnel/datasource/plugin/excel/ExcelAConfiguration.java
+5
-3
ExcelDataSourceFactory.java
...unnel/datasource/plugin/excel/ExcelDataSourceFactory.java
+4
-3
ExcelDatasourceChannel.java
...unnel/datasource/plugin/excel/ExcelDatasourceChannel.java
+8
-6
ExcelOptionRule.java
...he/seatunnel/datasource/plugin/excel/ExcelOptionRule.java
+0
-1
HttpAConfiguration.java
.../seatunnel/datasource/plugin/http/HttpAConfiguration.java
+2
-4
HttpClientService.java
...e/seatunnel/datasource/plugin/http/HttpClientService.java
+0
-6
HttpConfiguration.java
...e/seatunnel/datasource/plugin/http/HttpConfiguration.java
+3
-5
HttpDatasourceChannel.java
...atunnel/datasource/plugin/http/HttpDatasourceChannel.java
+19
-15
HttpOptionRule.java
...ache/seatunnel/datasource/plugin/http/HttpOptionRule.java
+1
-6
pom.xml
...atunnel-datasource-plugins/datasource-jdbc-access/pom.xml
+10
-0
AccessDataSourceConfig.java
...datasource/plugin/access/jdbc/AccessDataSourceConfig.java
+13
-6
AccessJdbcDataSourceChannel.java
...ource/plugin/access/jdbc/AccessJdbcDataSourceChannel.java
+0
-0
AccessOptionRule.java
...unnel/datasource/plugin/access/jdbc/AccessOptionRule.java
+6
-6
pom.xml
...atunnel-datasource-plugins/datasource-jdbc-dameng/pom.xml
+1
-1
DamengDataSourceConfig.java
...datasource/plugin/demeng/jdbc/DamengDataSourceConfig.java
+8
-7
DamengJdbcDataSourceChannel.java
...ource/plugin/demeng/jdbc/DamengJdbcDataSourceChannel.java
+0
-0
DamengJdbcDataSourceFactory.java
...ource/plugin/demeng/jdbc/DamengJdbcDataSourceFactory.java
+4
-4
DamengOptionRule.java
...unnel/datasource/plugin/demeng/jdbc/DamengOptionRule.java
+2
-4
pom.xml
...ource/seatunnel-datasource-plugins/datasource-xml/pom.xml
+1
-2
XMLAConfiguration.java
...he/seatunnel/datasource/plugin/xml/XMLAConfiguration.java
+4
-2
XMLDataSourceFactory.java
...seatunnel/datasource/plugin/xml/XMLDataSourceFactory.java
+3
-2
XMLDatasourceChannel.java
...seatunnel/datasource/plugin/xml/XMLDatasourceChannel.java
+8
-6
XMLOptionRule.java
...apache/seatunnel/datasource/plugin/xml/XMLOptionRule.java
+1
-2
pom.xml
seatunnel-datasource/seatunnel-datasource-plugins/pom.xml
+1
-1
没有找到文件。
seatunnel-datasource/seatunnel-datasource-plugins/datasource-csv/pom.xml
浏览文件 @
d18b47a1
...
...
@@ -141,7 +141,6 @@
<artifactId>
aws-java-sdk-bundle
</artifactId>
</dependency>
</dependencies>
</dependencies>
</project>
seatunnel-datasource/seatunnel-datasource-plugins/datasource-csv/src/main/java/org/apache/seatunnel/datasource/plugin/csv/CSVAConfiguration.java
浏览文件 @
d18b47a1
package
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
csv
;
import
lombok.extern.slf4j.Slf4j
;
import
org.apache.hadoop.conf.Configuration
;
import
org.apache.seatunnel.shade.com.typesafe.config.Config
;
import
org.apache.seatunnel.shade.com.typesafe.config.ConfigFactory
;
import
org.apache.hadoop.conf.Configuration
;
import
lombok.extern.slf4j.Slf4j
;
import
java.util.Map
;
@Slf4j
...
...
@@ -29,7 +31,7 @@ public class CSVAConfiguration {
throw
new
IllegalArgumentException
(
"S3 datasource endpoint is null, please check your config"
);
}
String
bucket
=
s3Options
.
get
(
CSVOptionRule
.
BUCKET
.
key
());
String
bucket
=
s3Options
.
get
(
CSVOptionRule
.
BUCKET
.
key
());
String
protocol
=
DEFAULT_PROTOCOL
;
if
(
bucket
.
startsWith
(
S3A_PROTOCOL
))
{
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-csv/src/main/java/org/apache/seatunnel/datasource/plugin/csv/CSVDataSourceFactory.java
浏览文件 @
d18b47a1
...
...
@@ -17,13 +17,14 @@
package
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
csv
;
import
com.google.auto.service.AutoService
;
import
com.google.common.collect.Sets
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceChannel
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceFactory
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourcePluginInfo
;
import
org.apache.seatunnel.datasource.plugin.api.DatasourcePluginTypeEnum
;
import
com.google.auto.service.AutoService
;
import
com.google.common.collect.Sets
;
import
java.util.Set
;
@AutoService
(
DataSourceFactory
.
class
)
...
...
@@ -52,6 +53,6 @@ public class CSVDataSourceFactory implements DataSourceFactory {
@Override
public
DataSourceChannel
createChannel
()
{
return
CSVDatasourceChannel
.
getInstance
();
return
CSVDatasourceChannel
.
getInstance
();
}
}
seatunnel-datasource/seatunnel-datasource-plugins/datasource-csv/src/main/java/org/apache/seatunnel/datasource/plugin/csv/CSVDatasourceChannel.java
浏览文件 @
d18b47a1
...
...
@@ -17,18 +17,20 @@
package
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
csv
;
import
org.apache.seatunnel.api.configuration.util.OptionRule
;
import
org.apache.seatunnel.common.utils.SeaTunnelException
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceChannel
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourcePluginException
;
import
org.apache.seatunnel.datasource.plugin.api.model.TableField
;
import
org.apache.commons.lang3.StringUtils
;
import
com.alibaba.fastjson2.JSON
;
import
io.minio.*
;
import
io.minio.errors.*
;
import
io.minio.messages.Bucket
;
import
io.minio.messages.Item
;
import
lombok.NonNull
;
import
org.apache.commons.lang3.StringUtils
;
import
org.apache.seatunnel.api.configuration.util.OptionRule
;
import
org.apache.seatunnel.common.utils.SeaTunnelException
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceChannel
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourcePluginException
;
import
org.apache.seatunnel.datasource.plugin.api.model.TableField
;
import
java.io.BufferedReader
;
import
java.io.IOException
;
...
...
@@ -293,7 +295,7 @@ public class CSVDatasourceChannel implements DataSourceChannel {
return
all
;
}
public
CSVClientService
createS3Client
(
Map
<
String
,
String
>
requestParams
)
{
public
CSVClientService
createS3Client
(
Map
<
String
,
String
>
requestParams
)
{
int
i
=
requestParams
.
get
(
"fs.s3a.endpoint"
).
lastIndexOf
(
":"
);
String
endpoint
=
requestParams
.
get
(
"fs.s3a.endpoint"
)
+
""
;
Integer
port
=
...
...
@@ -304,7 +306,7 @@ public class CSVDatasourceChannel implements DataSourceChannel {
String
password
=
requestParams
.
get
(
"secret_key"
)
+
""
;
// String bucket = requestParams.get("bucket") + "";
try
{
s3ClientService
=
new
CSVClientService
(
endpoint
,
provider
,
username
,
password
,
port
);
s3ClientService
=
new
CSVClientService
(
endpoint
,
provider
,
username
,
password
,
port
);
return
s3ClientService
;
}
catch
(
Exception
e
)
{
throw
new
SeaTunnelException
(
"创建Mqtt客户端错误!"
);
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-csv/src/main/java/org/apache/seatunnel/datasource/plugin/csv/CSVOptionRule.java
浏览文件 @
d18b47a1
...
...
@@ -21,7 +21,6 @@ import org.apache.seatunnel.api.configuration.Option;
import
org.apache.seatunnel.api.configuration.Options
;
import
org.apache.seatunnel.api.configuration.util.OptionRule
;
import
java.util.Arrays
;
import
java.util.Map
;
public
class
CSVOptionRule
{
...
...
@@ -129,7 +128,6 @@ public class CSVOptionRule {
.
build
();
}
public
enum
S3aAwsCredentialsProvider
{
SimpleAWSCredentialsProvider
(
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"
),
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-excel/pom.xml
浏览文件 @
d18b47a1
...
...
@@ -141,7 +141,6 @@
<artifactId>
aws-java-sdk-bundle
</artifactId>
</dependency>
</dependencies>
</dependencies>
</project>
seatunnel-datasource/seatunnel-datasource-plugins/datasource-excel/src/main/java/org/apache/seatunnel/datasource/plugin/excel/ExcelAConfiguration.java
浏览文件 @
d18b47a1
package
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
excel
;
import
lombok.extern.slf4j.Slf4j
;
import
org.apache.hadoop.conf.Configuration
;
import
org.apache.seatunnel.shade.com.typesafe.config.Config
;
import
org.apache.seatunnel.shade.com.typesafe.config.ConfigFactory
;
import
org.apache.hadoop.conf.Configuration
;
import
lombok.extern.slf4j.Slf4j
;
import
java.util.Map
;
@Slf4j
...
...
@@ -29,7 +31,7 @@ public class ExcelAConfiguration {
throw
new
IllegalArgumentException
(
"S3 datasource endpoint is null, please check your config"
);
}
String
bucket
=
s3Options
.
get
(
ExcelOptionRule
.
BUCKET
.
key
());
String
bucket
=
s3Options
.
get
(
ExcelOptionRule
.
BUCKET
.
key
());
String
protocol
=
DEFAULT_PROTOCOL
;
if
(
bucket
.
startsWith
(
S3A_PROTOCOL
))
{
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-excel/src/main/java/org/apache/seatunnel/datasource/plugin/excel/ExcelDataSourceFactory.java
浏览文件 @
d18b47a1
...
...
@@ -17,13 +17,14 @@
package
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
excel
;
import
com.google.auto.service.AutoService
;
import
com.google.common.collect.Sets
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceChannel
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceFactory
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourcePluginInfo
;
import
org.apache.seatunnel.datasource.plugin.api.DatasourcePluginTypeEnum
;
import
com.google.auto.service.AutoService
;
import
com.google.common.collect.Sets
;
import
java.util.Set
;
@AutoService
(
DataSourceFactory
.
class
)
...
...
@@ -52,6 +53,6 @@ public class ExcelDataSourceFactory implements DataSourceFactory {
@Override
public
DataSourceChannel
createChannel
()
{
return
ExcelDatasourceChannel
.
getInstance
();
return
ExcelDatasourceChannel
.
getInstance
();
}
}
seatunnel-datasource/seatunnel-datasource-plugins/datasource-excel/src/main/java/org/apache/seatunnel/datasource/plugin/excel/ExcelDatasourceChannel.java
浏览文件 @
d18b47a1
...
...
@@ -17,18 +17,20 @@
package
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
excel
;
import
org.apache.seatunnel.api.configuration.util.OptionRule
;
import
org.apache.seatunnel.common.utils.SeaTunnelException
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceChannel
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourcePluginException
;
import
org.apache.seatunnel.datasource.plugin.api.model.TableField
;
import
org.apache.commons.lang3.StringUtils
;
import
com.alibaba.fastjson2.JSON
;
import
io.minio.*
;
import
io.minio.errors.*
;
import
io.minio.messages.Bucket
;
import
io.minio.messages.Item
;
import
lombok.NonNull
;
import
org.apache.commons.lang3.StringUtils
;
import
org.apache.seatunnel.api.configuration.util.OptionRule
;
import
org.apache.seatunnel.common.utils.SeaTunnelException
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceChannel
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourcePluginException
;
import
org.apache.seatunnel.datasource.plugin.api.model.TableField
;
import
java.io.BufferedReader
;
import
java.io.IOException
;
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-excel/src/main/java/org/apache/seatunnel/datasource/plugin/excel/ExcelOptionRule.java
浏览文件 @
d18b47a1
...
...
@@ -129,7 +129,6 @@ public class ExcelOptionRule {
.
build
();
}
public
enum
S3aAwsCredentialsProvider
{
SimpleAWSCredentialsProvider
(
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"
),
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-http/src/main/java/org/apache/seatunnel/datasource/plugin/http/HttpAConfiguration.java
浏览文件 @
d18b47a1
...
...
@@ -27,17 +27,15 @@ public class HttpAConfiguration {
public
static
HttpConfiguration
getConfiguration
(
Map
<
String
,
String
>
ftpOption
)
{
if
(!
ftpOption
.
containsKey
(
HttpOptionRule
.
URL
.
key
()))
{
throw
new
IllegalArgumentException
(
"url is null, please check your config"
);
throw
new
IllegalArgumentException
(
"url is null, please check your config"
);
}
HttpConfiguration
httpAConfiguration
=
new
HttpConfiguration
();
HttpConfiguration
httpAConfiguration
=
new
HttpConfiguration
();
httpAConfiguration
.
setUrl
(
HttpOptionRule
.
URL
.
key
());
httpAConfiguration
.
setMethod
(
HttpOptionRule
.
METHOD
.
key
());
httpAConfiguration
.
setToken
(
HttpOptionRule
.
TOKEN
.
key
());
httpAConfiguration
.
setRequest_params
(
HttpOptionRule
.
REQUEST_PARAMS
.
key
());
return
httpAConfiguration
;
}
}
seatunnel-datasource/seatunnel-datasource-plugins/datasource-http/src/main/java/org/apache/seatunnel/datasource/plugin/http/HttpClientService.java
浏览文件 @
d18b47a1
package
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
http
;
import
org.apache.http.client.HttpClient
;
import
org.apache.http.impl.client.HttpClients
;
...
...
@@ -8,14 +7,9 @@ public class HttpClientService {
public
static
HttpClient
connect
(
HttpConfiguration
conf
)
throws
Exception
{
// 创建HttpClient实例
HttpClient
client
=
HttpClients
.
createDefault
();
return
client
;
}
}
seatunnel-datasource/seatunnel-datasource-plugins/datasource-http/src/main/java/org/apache/seatunnel/datasource/plugin/http/HttpConfiguration.java
浏览文件 @
d18b47a1
...
...
@@ -10,12 +10,11 @@ public class HttpConfiguration {
public
HttpConfiguration
()
{}
public
HttpConfiguration
(
String
url
,
String
method
,
String
token
,
String
request_params
)
{
public
HttpConfiguration
(
String
url
,
String
method
,
String
token
,
String
request_params
)
{
this
.
url
=
url
;
this
.
token
=
token
;
this
.
method
=
method
;
this
.
request_params
=
request_params
;
this
.
method
=
method
;
this
.
request_params
=
request_params
;
}
public
String
getUrl
()
{
...
...
@@ -26,7 +25,6 @@ public class HttpConfiguration {
this
.
url
=
url
;
}
public
String
getMethod
()
{
return
method
;
}
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-http/src/main/java/org/apache/seatunnel/datasource/plugin/http/HttpDatasourceChannel.java
浏览文件 @
d18b47a1
...
...
@@ -17,23 +17,21 @@
package
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
http
;
import
org.apache.seatunnel.api.configuration.util.OptionRule
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceChannel
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourcePluginException
;
import
org.apache.seatunnel.datasource.plugin.api.model.TableField
;
import
org.apache.commons.lang3.StringUtils
;
import
org.apache.http.HttpResponse
;
import
org.apache.http.client.HttpClient
;
import
org.apache.http.client.methods.*
;
import
org.apache.http.entity.StringEntity
;
import
org.apache.http.util.EntityUtils
;
import
org.apache.seatunnel.api.configuration.util.OptionRule
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceChannel
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourcePluginException
;
import
org.apache.seatunnel.datasource.plugin.api.model.TableField
;
import
lombok.NonNull
;
import
java.net.URI
;
import
java.util.List
;
import
java.util.Map
;
import
java.util.Objects
;
...
...
@@ -75,7 +73,8 @@ public class HttpDatasourceChannel implements DataSourceChannel {
System
.
out
.
println
(
"url:"
+
url
);
HttpGet
httpGet
=
new
HttpGet
(
url
);
if
(
StringUtils
.
isNotBlank
(
token
))
{
httpGet
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
httpGet
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
}
// 执行请求并获得响应
...
...
@@ -83,7 +82,8 @@ public class HttpDatasourceChannel implements DataSourceChannel {
}
else
if
(
StringUtils
.
isBlank
(
method
)
||
"POST"
.
equals
(
method
.
toUpperCase
()))
{
HttpPost
httpPost
=
new
HttpPost
(
url
);
if
(
StringUtils
.
isNotBlank
(
token
))
{
httpPost
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
httpPost
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
}
// 设置请求体(例如:JSON数据)
StringEntity
requestEntity
=
new
StringEntity
(
parmams
);
...
...
@@ -94,7 +94,8 @@ public class HttpDatasourceChannel implements DataSourceChannel {
}
else
if
(
StringUtils
.
isBlank
(
method
)
||
"PUT"
.
equals
(
method
.
toUpperCase
()))
{
HttpPut
httpPut
=
new
HttpPut
(
url
);
if
(
StringUtils
.
isNotBlank
(
token
))
{
httpPut
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
httpPut
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
}
// 设置请求体(例如:JSON数据)
StringEntity
requestEntity
=
new
StringEntity
(
parmams
);
...
...
@@ -110,7 +111,8 @@ public class HttpDatasourceChannel implements DataSourceChannel {
HttpDelete
httpDelete
=
new
HttpDelete
(
url
);
if
(
StringUtils
.
isNotBlank
(
token
))
{
httpDelete
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
httpDelete
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
}
// 执行请求并获得响应
...
...
@@ -119,7 +121,8 @@ public class HttpDatasourceChannel implements DataSourceChannel {
}
else
if
(
StringUtils
.
isBlank
(
method
)
||
"PATCH"
.
equals
(
method
.
toUpperCase
()))
{
HttpPatch
httpPatch
=
new
HttpPatch
(
url
);
if
(
StringUtils
.
isNotBlank
(
token
))
{
httpPatch
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
httpPatch
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
}
// 设置请求体(例如:JSON数据)
StringEntity
requestEntity
=
new
StringEntity
(
parmams
);
...
...
@@ -133,7 +136,8 @@ public class HttpDatasourceChannel implements DataSourceChannel {
}
HttpOptions
httpOptions
=
new
HttpOptions
(
url
);
if
(
StringUtils
.
isNotBlank
(
token
))
{
httpOptions
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
httpOptions
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
}
// 执行请求并获得响应
...
...
@@ -145,7 +149,8 @@ public class HttpDatasourceChannel implements DataSourceChannel {
}
HttpHead
httpHead
=
new
HttpHead
(
url
);
if
(
StringUtils
.
isNotBlank
(
token
))
{
httpHead
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
httpHead
.
setHeader
(
"Authorization"
,
"Bearer "
+
token
.
replace
(
"Bearer "
,
""
).
trim
());
}
// 执行请求并获得响应
...
...
@@ -160,7 +165,6 @@ public class HttpDatasourceChannel implements DataSourceChannel {
String
responseBody
=
EntityUtils
.
toString
(
response
.
getEntity
());
System
.
out
.
println
(
"Response Body: "
+
responseBody
);
if
(
statusCode
==
200
)
{
return
true
;
}
else
{
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-http/src/main/java/org/apache/seatunnel/datasource/plugin/http/HttpOptionRule.java
浏览文件 @
d18b47a1
...
...
@@ -47,11 +47,8 @@ public class HttpOptionRule {
.
noDefaultValue
()
.
withDescription
(
"the http user token to use when connecting to the broker"
);
public
static
OptionRule
optionRule
()
{
return
OptionRule
.
builder
().
required
(
URL
,
METHOD
).
optional
(
TOKEN
,
REQUEST_PARAMS
).
build
();
return
OptionRule
.
builder
().
required
(
URL
,
METHOD
).
optional
(
TOKEN
,
REQUEST_PARAMS
).
build
();
}
public
static
OptionRule
metadataRule
()
{
...
...
@@ -59,8 +56,6 @@ public class HttpOptionRule {
}
public
enum
FileFormat
{
JSON
(
"json"
),
;
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-access/pom.xml
浏览文件 @
d18b47a1
...
...
@@ -56,7 +56,17 @@
<artifactId>
ucanaccess
</artifactId>
<version>
5.0.1
</version>
</dependency>
<dependency>
<groupId>
com.cdzhiyong.unified
</groupId>
<artifactId>
minio-spring-boot-starter
</artifactId>
<version>
1.0.0
</version>
</dependency>
<dependency>
<groupId>
org.apache.httpcomponents
</groupId>
<artifactId>
httpclient
</artifactId>
<version>
4.5.14
</version>
</dependency>
</dependencies>
</project>
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-access/src/main/java/org/apache/seatunnel/datasource/plugin/access/jdbc/AccessDataSourceConfig.java
浏览文件 @
d18b47a1
...
...
@@ -37,15 +37,22 @@ public class AccessDataSourceConfig {
.
type
(
DatasourcePluginTypeEnum
.
DATABASE
.
getCode
())
.
build
();
public
static
final
Set
<
String
>
MYSQL_SYSTEM_DATABASES
=
Sets
.
newHashSet
(
"SYSTEM"
,
"ROLL"
);
public
static
final
Set
<
String
>
MYSQL_SYSTEM_DATABASES
=
Sets
.
newHashSet
(
"SYSTEM"
,
"ROLL"
);
public
static
final
OptionRule
OPTION_RULE
=
OptionRule
.
builder
()
.
required
(
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
access
.
jdbc
.
AccessOptionRule
.
URL
,
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
access
.
jdbc
.
AccessOptionRule
.
DRIVER
)
.
optional
(
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
access
.
jdbc
.
AccessOptionRule
.
USER
,
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
access
.
jdbc
.
AccessOptionRule
.
PASSWORD
)
.
required
(
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
access
.
jdbc
.
AccessOptionRule
.
URL
,
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
access
.
jdbc
.
AccessOptionRule
.
DRIVER
)
.
optional
(
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
access
.
jdbc
.
AccessOptionRule
.
USER
,
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
access
.
jdbc
.
AccessOptionRule
.
PASSWORD
)
.
build
();
// public static final OptionRule METADATA_RULE =
// OptionRule.builder().required(org.apache.seatunnel.datasource.plugin.demeng.jdbc.DemengOptionRule.DATABASE, org.apache.seatunnel.datasource.plugin.demeng.jdbc.DemengOptionRule.TABLE).build();
// public static final OptionRule METADATA_RULE =
//
// OptionRule.builder().required(org.apache.seatunnel.datasource.plugin.demeng.jdbc.DemengOptionRule.DATABASE, org.apache.seatunnel.datasource.plugin.demeng.jdbc.DemengOptionRule.TABLE).build();
}
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-access/src/main/java/org/apache/seatunnel/datasource/plugin/access/jdbc/AccessJdbcDataSourceChannel.java
浏览文件 @
d18b47a1
差异被折叠。
点击展开。
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-access/src/main/java/org/apache/seatunnel/datasource/plugin/access/jdbc/AccessOptionRule.java
浏览文件 @
d18b47a1
...
...
@@ -27,8 +27,7 @@ public class AccessOptionRule {
.
stringType
()
.
noDefaultValue
()
.
withDescription
(
"jdbc url, eg:"
+
" http://localhost:9000/bucket/filename.mdb"
);
"jdbc url, eg:"
+
" http://localhost:9000/bucket/filename.mdb"
);
public
static
final
Option
<
String
>
USER
=
Options
.
key
(
"user"
).
stringType
().
noDefaultValue
().
withDescription
(
"jdbc user"
);
...
...
@@ -36,11 +35,12 @@ public class AccessOptionRule {
public
static
final
Option
<
String
>
PASSWORD
=
Options
.
key
(
"password"
).
stringType
().
noDefaultValue
().
withDescription
(
"jdbc password"
);
// public static final Option<String> DATABASE =
// Options.key("database").stringType().noDefaultValue().withDescription("jdbc database");
// public static final Option<String> DATABASE =
// Options.key("database").stringType().noDefaultValue().withDescription("jdbc
// database");
// public static final Option<String> TABLE =
// Options.key("table").stringType().noDefaultValue().withDescription("jdbc table");
// public static final Option<String> TABLE =
// Options.key("table").stringType().noDefaultValue().withDescription("jdbc table");
public
static
final
Option
<
DriverType
>
DRIVER
=
Options
.
key
(
"driver"
)
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-d
e
meng/pom.xml
→
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-d
a
meng/pom.xml
浏览文件 @
d18b47a1
...
...
@@ -22,7 +22,7 @@
<version>
1.0.0-SNAPSHOT
</version>
</parent>
<artifactId>
datasource-jdbc-d
e
meng
</artifactId>
<artifactId>
datasource-jdbc-d
a
meng
</artifactId>
<properties>
<mysql-connector.version>
8.0.28
</mysql-connector.version>
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-d
emeng/src/main/java/org/apache/seatunnel/datasource/plugin/demeng/jdbc/De
mengDataSourceConfig.java
→
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-d
ameng/src/main/java/org/apache/seatunnel/datasource/plugin/demeng/jdbc/Da
mengDataSourceConfig.java
浏览文件 @
d18b47a1
...
...
@@ -25,9 +25,9 @@ import com.google.common.collect.Sets;
import
java.util.Set
;
public
class
D
e
mengDataSourceConfig
{
public
class
D
a
mengDataSourceConfig
{
public
static
final
String
PLUGIN_NAME
=
"JDBC-D
e
meng"
;
public
static
final
String
PLUGIN_NAME
=
"JDBC-D
a
meng"
;
public
static
final
DataSourcePluginInfo
MYSQL_DATASOURCE_PLUGIN_INFO
=
DataSourcePluginInfo
.
builder
()
...
...
@@ -37,15 +37,16 @@ public class DemengDataSourceConfig {
.
type
(
DatasourcePluginTypeEnum
.
DATABASE
.
getCode
())
.
build
();
public
static
final
Set
<
String
>
MYSQL_SYSTEM_DATABASES
=
Sets
.
newHashSet
(
"SYSTEM"
,
"ROLL"
);
public
static
final
Set
<
String
>
MYSQL_SYSTEM_DATABASES
=
Sets
.
newHashSet
(
"SYSTEM"
,
"ROLL"
);
public
static
final
OptionRule
OPTION_RULE
=
OptionRule
.
builder
()
.
required
(
D
emengOptionRule
.
URL
,
De
mengOptionRule
.
DRIVER
)
.
optional
(
D
emengOptionRule
.
USER
,
De
mengOptionRule
.
PASSWORD
)
.
required
(
D
amengOptionRule
.
URL
,
Da
mengOptionRule
.
DRIVER
)
.
optional
(
D
amengOptionRule
.
USER
,
Da
mengOptionRule
.
PASSWORD
)
.
build
();
public
static
final
OptionRule
METADATA_RULE
=
OptionRule
.
builder
().
required
(
DemengOptionRule
.
DATABASE
,
DemengOptionRule
.
TABLE
).
build
();
OptionRule
.
builder
()
.
required
(
DamengOptionRule
.
DATABASE
,
DamengOptionRule
.
TABLE
)
.
build
();
}
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-d
emeng/src/main/java/org/apache/seatunnel/datasource/plugin/demeng/jdbc/De
mengJdbcDataSourceChannel.java
→
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-d
ameng/src/main/java/org/apache/seatunnel/datasource/plugin/demeng/jdbc/Da
mengJdbcDataSourceChannel.java
浏览文件 @
d18b47a1
差异被折叠。
点击展开。
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-d
emeng/src/main/java/org/apache/seatunnel/datasource/plugin/demeng/jdbc/De
mengJdbcDataSourceFactory.java
→
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-d
ameng/src/main/java/org/apache/seatunnel/datasource/plugin/demeng/jdbc/Da
mengJdbcDataSourceFactory.java
浏览文件 @
d18b47a1
...
...
@@ -29,20 +29,20 @@ import java.util.Set;
@Slf4j
@AutoService
(
DataSourceFactory
.
class
)
public
class
D
e
mengJdbcDataSourceFactory
implements
DataSourceFactory
{
public
class
D
a
mengJdbcDataSourceFactory
implements
DataSourceFactory
{
@Override
public
String
factoryIdentifier
()
{
return
D
e
mengDataSourceConfig
.
PLUGIN_NAME
;
return
D
a
mengDataSourceConfig
.
PLUGIN_NAME
;
}
@Override
public
Set
<
DataSourcePluginInfo
>
supportedDataSources
()
{
return
Sets
.
newHashSet
(
D
e
mengDataSourceConfig
.
MYSQL_DATASOURCE_PLUGIN_INFO
);
return
Sets
.
newHashSet
(
D
a
mengDataSourceConfig
.
MYSQL_DATASOURCE_PLUGIN_INFO
);
}
@Override
public
DataSourceChannel
createChannel
()
{
return
new
D
e
mengJdbcDataSourceChannel
();
return
new
D
a
mengJdbcDataSourceChannel
();
}
}
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-d
emeng/src/main/java/org/apache/seatunnel/datasource/plugin/demeng/jdbc/De
mengOptionRule.java
→
seatunnel-datasource/seatunnel-datasource-plugins/datasource-jdbc-d
ameng/src/main/java/org/apache/seatunnel/datasource/plugin/demeng/jdbc/Da
mengOptionRule.java
浏览文件 @
d18b47a1
...
...
@@ -20,15 +20,13 @@ package org.apache.seatunnel.datasource.plugin.demeng.jdbc;
import
org.apache.seatunnel.api.configuration.Option
;
import
org.apache.seatunnel.api.configuration.Options
;
public
class
D
e
mengOptionRule
{
public
class
D
a
mengOptionRule
{
public
static
final
Option
<
String
>
URL
=
Options
.
key
(
"url"
)
.
stringType
()
.
noDefaultValue
()
.
withDescription
(
"jdbc url, eg:"
+
" jdbc:dm://localhost:5236"
);
.
withDescription
(
"jdbc url, eg:"
+
" jdbc:dm://localhost:5236"
);
public
static
final
Option
<
String
>
USER
=
Options
.
key
(
"user"
).
stringType
().
noDefaultValue
().
withDescription
(
"jdbc user"
);
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-xml/pom.xml
浏览文件 @
d18b47a1
...
...
@@ -141,7 +141,6 @@
<artifactId>
aws-java-sdk-bundle
</artifactId>
</dependency>
</dependencies>
</dependencies>
</project>
seatunnel-datasource/seatunnel-datasource-plugins/datasource-xml/src/main/java/org/apache/seatunnel/datasource/plugin/xml/XMLAConfiguration.java
浏览文件 @
d18b47a1
package
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
xml
;
import
lombok.extern.slf4j.Slf4j
;
import
org.apache.hadoop.conf.Configuration
;
import
org.apache.seatunnel.shade.com.typesafe.config.Config
;
import
org.apache.seatunnel.shade.com.typesafe.config.ConfigFactory
;
import
org.apache.hadoop.conf.Configuration
;
import
lombok.extern.slf4j.Slf4j
;
import
java.util.Map
;
@Slf4j
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-xml/src/main/java/org/apache/seatunnel/datasource/plugin/xml/XMLDataSourceFactory.java
浏览文件 @
d18b47a1
...
...
@@ -17,13 +17,14 @@
package
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
xml
;
import
com.google.auto.service.AutoService
;
import
com.google.common.collect.Sets
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceChannel
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceFactory
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourcePluginInfo
;
import
org.apache.seatunnel.datasource.plugin.api.DatasourcePluginTypeEnum
;
import
com.google.auto.service.AutoService
;
import
com.google.common.collect.Sets
;
import
java.util.Set
;
@AutoService
(
DataSourceFactory
.
class
)
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-xml/src/main/java/org/apache/seatunnel/datasource/plugin/xml/XMLDatasourceChannel.java
浏览文件 @
d18b47a1
...
...
@@ -17,18 +17,20 @@
package
org
.
apache
.
seatunnel
.
datasource
.
plugin
.
xml
;
import
org.apache.seatunnel.api.configuration.util.OptionRule
;
import
org.apache.seatunnel.common.utils.SeaTunnelException
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceChannel
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourcePluginException
;
import
org.apache.seatunnel.datasource.plugin.api.model.TableField
;
import
org.apache.commons.lang3.StringUtils
;
import
com.alibaba.fastjson2.JSON
;
import
io.minio.*
;
import
io.minio.errors.*
;
import
io.minio.messages.Bucket
;
import
io.minio.messages.Item
;
import
lombok.NonNull
;
import
org.apache.commons.lang3.StringUtils
;
import
org.apache.seatunnel.api.configuration.util.OptionRule
;
import
org.apache.seatunnel.common.utils.SeaTunnelException
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourceChannel
;
import
org.apache.seatunnel.datasource.plugin.api.DataSourcePluginException
;
import
org.apache.seatunnel.datasource.plugin.api.model.TableField
;
import
java.io.BufferedReader
;
import
java.io.IOException
;
...
...
seatunnel-datasource/seatunnel-datasource-plugins/datasource-xml/src/main/java/org/apache/seatunnel/datasource/plugin/xml/XMLOptionRule.java
浏览文件 @
d18b47a1
...
...
@@ -21,7 +21,6 @@ import org.apache.seatunnel.api.configuration.Option;
import
org.apache.seatunnel.api.configuration.Options
;
import
org.apache.seatunnel.api.configuration.util.OptionRule
;
import
java.util.Arrays
;
import
java.util.Map
;
public
class
XMLOptionRule
{
...
...
@@ -121,7 +120,7 @@ public class XMLOptionRule {
return
OptionRule
.
builder
()
.
required
(
PATH
,
TYPE
)
.
conditional
(
TYPE
,
FileFormat
.
XML
.
type
,
DELIMITER
)
.
conditional
(
TYPE
,
FileFormat
.
XML
.
type
,
SCHEMA
)
.
conditional
(
TYPE
,
FileFormat
.
XML
.
type
,
SCHEMA
)
.
optional
(
PARSE_PARSE_PARTITION_FROM_PATH
)
.
optional
(
DATE_FORMAT
)
.
optional
(
DATETIME_FORMAT
)
...
...
seatunnel-datasource/seatunnel-datasource-plugins/pom.xml
浏览文件 @
d18b47a1
...
...
@@ -50,7 +50,7 @@
<module>
datasource-redis
</module>
<module>
datasource-rabbitmq
</module>
<module>
datasource-ftp
</module>
<module>
datasource-jdbc-d
e
meng
</module>
<module>
datasource-jdbc-d
a
meng
</module>
<module>
datasource-jdbc-access
</module>
<module>
datasource-http
</module>
<module>
datasource-xml
</module>
...
...
编写
预览
Markdown
格式
0%
重试
或
添加新文件
添加附件
取消
您添加了
0
人
到此讨论。请谨慎行事。
请先完成此评论的编辑!
取消
请
注册
或者
登录
后发表评论