Skip to content

Commit 1569e37

Browse files
committed
add fileSystem catalog supporting for tis, issue:datavane/tis#490
1 parent e52c4c0 commit 1569e37

67 files changed

Lines changed: 5716 additions & 88 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

pom.xml

Lines changed: 30 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -75,25 +75,25 @@
7575
<module>tis-local-dump-build</module>-->
7676
</modules>
7777
<profiles>
78-
<profile>
79-
<id>default-emr</id>
80-
<properties>
81-
<spark2.version>2.4.4</spark2.version>
82-
<hadoop-version>${hadoop2x-version}</hadoop-version>
83-
<spark.dist.dir.name>spark-${spark2.version}-bin-hadoop2.7</spark.dist.dir.name>
84-
<hudi.version>0.14.1</hudi.version>
85-
<hive.version>2.3.1</hive.version>
86-
<scala.binary.version>2.12</scala.binary.version>
87-
<scala.version>2.12.9</scala.version>
88-
<appname>all</appname>
89-
90-
<hudi>${hudi.version}</hudi>
91-
<spark>${spark2.version}</spark>
92-
<hive>${hive.version}</hive>
93-
<hadoop>${hadoop-version}</hadoop>
94-
95-
</properties>
96-
</profile>
78+
<!-- <profile>-->
79+
<!-- <id>default-emr</id>-->
80+
<!-- <properties>-->
81+
<!-- <spark2.version>2.4.4</spark2.version>-->
82+
<!-- <hadoop-version>${hadoop2x-version}</hadoop-version>-->
83+
<!-- <spark.dist.dir.name>spark-${spark2.version}-bin-hadoop2.7</spark.dist.dir.name>-->
84+
<!-- <hudi.version>0.14.1</hudi.version>-->
85+
<!-- <hive.version>2.3.1</hive.version>-->
86+
<!-- <scala.binary.version>2.12</scala.binary.version>-->
87+
<!-- <scala.version>2.12.9</scala.version>-->
88+
<!-- <appname>all</appname>-->
89+
90+
<!-- <hudi>${hudi.version}</hudi>-->
91+
<!-- <spark>${spark2.version}</spark>-->
92+
<!-- <hive>${hive.version}</hive>-->
93+
<!-- <hadoop>${hadoop-version}</hadoop>-->
94+
95+
<!-- </properties>-->
96+
<!-- </profile>-->
9797
<profile>
9898
<!--https://emr-next.console.aliyun.com/#/resource/all/create/ecs-->
9999
<id>aliyun-emr</id>
@@ -164,8 +164,10 @@
164164
<hudi.version>0.14.1</hudi.version>
165165
<!-- <revision>4.2.0</revision>-->
166166
<spark.dist.dir.name>spark-${spark2.version}-bin-hadoop2.7</spark.dist.dir.name>
167-
<hadoop2x-version>2.7.3</hadoop2x-version>
168-
<hadoop-version>${hadoop2x-version}</hadoop-version>
167+
<!-- <hadoop2x-version>2.7.3</hadoop2x-version>-->
168+
169+
<hadoop-version>${hadoop3-version}</hadoop-version>
170+
<!-- <hadoop-version>${hadoop2x-version}</hadoop-version>-->
169171
<hive.version>2.3.1</hive.version>
170172

171173
<maven-tpi-plugin.classpathDependentExcludes>
@@ -686,13 +688,13 @@
686688
</dependency>
687689

688690
<!-- 解决IDE对jdk.tools依赖的报错,该依赖实际上会被exclusion排除 -->
689-
<dependency>
690-
<groupId>jdk.tools</groupId>
691-
<artifactId>jdk.tools</artifactId>
692-
<version>1.8</version>
693-
<scope>system</scope>
694-
<systemPath>${java.home}/../lib/tools.jar</systemPath>
695-
</dependency>
691+
<!-- <dependency>-->
692+
<!-- <groupId>jdk.tools</groupId>-->
693+
<!-- <artifactId>jdk.tools</artifactId>-->
694+
<!-- <version>1.8</version>-->
695+
<!-- <scope>system</scope>-->
696+
<!-- <systemPath>${java.home}/../lib/tools.jar</systemPath>-->
697+
<!-- </dependency>-->
696698

697699
</dependencies>
698700
</dependencyManagement>

tis-datax/executor/tis-datax-executor/src/main/java/com/qlangtech/tis/datax/executor/BasicTISTableDumpProcessor.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
import com.qlangtech.tis.datax.DataxPrePostConsumer;
2626
import com.qlangtech.tis.datax.IDataXBatchPost;
2727
import com.qlangtech.tis.datax.IDataxProcessor;
28+
import com.qlangtech.tis.datax.IDataxReader;
2829
import com.qlangtech.tis.datax.IDataxWriter;
2930
import com.qlangtech.tis.datax.LifeCycleHook;
3031
import com.qlangtech.tis.datax.RpcUtils;
@@ -291,8 +292,10 @@ public static IRemoteTaskTrigger createDataXJob(AbstractExecContext execContext
291292
IDataxProcessor processor = execContext.getProcessor();
292293

293294
if (TisAppLaunch.isTestMock()) {
295+
IDataxReader reader = processor.getReader(null);
296+
ISelectedTab tab = reader.getSelectedTab(tableName);
294297
IDataXBatchPost dataXBatchPost = getDataXBatchPost(processor);
295-
DefaultTab tab = new DefaultTab(tableName);
298+
// DefaultTab tab = new DefaultTab(tableName);
296299
EntityName entityName = dataXBatchPost.parseEntity(tab);
297300
LifeCycleHook cycleHook = lifeCycleHookInfo.getRight();
298301
if (cycleHook == LifeCycleHook.Post) {

tis-datax/tis-aliyun-jindo-sdk-extends/pom.xml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,7 @@
3131
<module>tis-aliyun-jindo-sdk-extends-api</module>
3232
<module>tis-aliyun-jindo-sdk-extends-impl</module>
3333
<module>tis-datax-hdfs-aliyun-emr-plugin</module>
34-
<module>tis-aliyun-jindo-sdk-extends-hadoop2x</module>
34+
<!-- <module>tis-aliyun-jindo-sdk-extends-hadoop2x</module>-->
3535
<module>tis-aliyun-jindo-sdk-extends-hadoop3x</module>
3636
</modules>
3737

tis-datax/tis-datax-hdfs-plugin/src/main/java/com/qlangtech/tis/hdfs/impl/HdfsFileSystemFactory.java

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@
3737
import com.qlangtech.tis.manage.common.TisUTF8;
3838
import com.qlangtech.tis.offline.FileSystemFactory;
3939
import com.qlangtech.tis.plugin.IEndTypeGetter;
40+
import com.qlangtech.tis.plugin.amazon.s3.ReplayConfiguration;
4041
import com.qlangtech.tis.plugin.annotation.FormField;
4142
import com.qlangtech.tis.plugin.annotation.FormFieldType;
4243
import com.qlangtech.tis.plugin.annotation.Validator;
@@ -205,7 +206,7 @@ public static List<? extends Descriptor> filter(List<? extends Descriptor> descs
205206
}).collect(Collectors.toList());
206207
}
207208

208-
private static Configuration getConfiguration(String hdfsSiteContent, Boolean userHostname, Consumer<Configuration> cfgProcess) {
209+
private static ReplayConfiguration getConfiguration(String hdfsSiteContent, Boolean userHostname, Consumer<ReplayConfiguration> cfgProcess) {
209210

210211

211212
// final ClassLoader contextClassLoader = Thread.currentThread().getContextClassLoader();
@@ -221,7 +222,7 @@ private static Configuration getConfiguration(String hdfsSiteContent, Boolean us
221222

222223
try {
223224
return ClassloaderUtils.processByResetThreadClassloader(HdfsFileSystemFactory.class, () -> {
224-
Configuration conf = new Configuration();
225+
ReplayConfiguration conf = new ReplayConfiguration();
225226
try (InputStream input = new ByteArrayInputStream(hdfsSiteContent.getBytes(TisUTF8.get()))) {
226227
conf.addResource(input);
227228
}
@@ -259,7 +260,7 @@ private static Configuration getConfiguration(String hdfsSiteContent, Boolean us
259260
}
260261

261262
@Override
262-
public Configuration getConfiguration() {
263+
public ReplayConfiguration getConfiguration() {
263264
return getConfiguration(this.hdfsSiteContent, this.userHostname, (conf) -> setConfiguration(conf));
264265
}
265266

tis-datax/tis-datax-hdfs-plugin/src/main/java/com/qlangtech/tis/plugin/amazon/s3/AmazonS3FileSystemFactory.java

Lines changed: 31 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -13,11 +13,13 @@
1313
import com.qlangtech.tis.offline.FileSystemFactory;
1414
import com.qlangtech.tis.plugin.IEndTypeGetter;
1515
import com.qlangtech.tis.plugin.annotation.FormField;
16+
import com.qlangtech.tis.plugin.annotation.FormFieldType;
1617
import com.qlangtech.tis.plugin.annotation.Validator;
1718
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
1819
import com.qlangtech.tis.util.ClassloaderUtils;
1920
import org.apache.commons.lang.StringUtils;
2021
import org.apache.hadoop.conf.Configuration;
22+
import org.apache.hadoop.fs.s3a.Constants;
2123
import org.apache.hadoop.hdfs.DFSConfigKeys;
2224
import org.slf4j.Logger;
2325
import org.slf4j.LoggerFactory;
@@ -42,21 +44,24 @@ public class AmazonS3FileSystemFactory extends FileSystemFactory implements ITIS
4244
@FormField(ordinal = 1, validate = {Validator.require, Validator.url})
4345
public String endpoint;
4446

45-
@FormField(ordinal = 2, validate = {Validator.require, Validator.absolute_path})
47+
/**
48+
* 是存储桶名称,相当于一个顶级容器
49+
*/
50+
@FormField(ordinal = 2, validate = {Validator.require, Validator.identity})
51+
public String bucket;
52+
53+
@FormField(ordinal = 3, validate = {Validator.require, Validator.absolute_path})
4654
public String rootDir;
4755

4856
//<!-- 可选:如果遇到 region 相关报错,可以设置一个默认 region -->
49-
@FormField(ordinal = 3, validate = {Validator.identity})
57+
@FormField(ordinal = 4, validate = {Validator.identity})
5058
public String region;
5159

5260
@FormField(ordinal = 5, validate = {Validator.require})
5361
public UserToken userToken;
5462

55-
/**
56-
* 是存储桶名称,相当于一个顶级容器
57-
*/
58-
@FormField(ordinal = 7, validate = {Validator.require, Validator.identity})
59-
public String bucket;
63+
@FormField(ordinal = 6, type = FormFieldType.ENUM, validate = {Validator.require})
64+
public Boolean pathStyleAccess ;
6065

6166

6267
// public String accessKey;
@@ -68,10 +73,10 @@ public class AmazonS3FileSystemFactory extends FileSystemFactory implements ITIS
6873
private transient ITISFileSystem fileSystem;
6974

7075
@Override
71-
public Configuration getConfiguration() {
76+
public ReplayConfiguration getConfiguration() {
7277
try {
7378
return ClassloaderUtils.processByResetThreadClassloader(AmazonS3FileSystemFactory.class, () -> {
74-
final Configuration conf = new Configuration();
79+
final ReplayConfiguration conf = new ReplayConfiguration();
7580
// try (InputStream input = new ByteArrayInputStream(hdfsSiteContent.getBytes(TisUTF8.get()))) {
7681
// conf.addResource(input);
7782
// }
@@ -80,7 +85,15 @@ public Configuration getConfiguration() {
8085
// conf.setBoolean(DFSConfigKeys.DFS_CLIENT_RETRY_POLICY_ENABLED_KEY, false);
8186
// conf.set(FsPermission.UMASK_LABEL, "000");
8287
// fs.defaultFS
83-
conf.set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem");
88+
conf.set("fs.s3a.impl", org.apache.hadoop.fs.s3a.S3AFileSystem.class.getName());
89+
90+
/**
91+
* 显式注册 LocalFileSystem,防止在 TIS 插件 ClassLoader 隔离环境下
92+
* S3AFileSystem 创建本地临时文件时找不到 "file" scheme 的实现
93+
* S3AFileSystem 需要本地临时文件:当 Paimon 通过 HadoopFileIO 写入 S3/MinIO 时,S3AFileSystem 内部使用 DiskBlockFactory 先在本地磁盘创建临时文件缓冲数据,然后再上传。创建临时文件时调用了 FileSystem.getLocal(conf),这需要查找 "file"
94+
* scheme 对应的 LocalFileSystem 实现。
95+
*/
96+
conf.set("fs.file.impl", org.apache.hadoop.fs.LocalFileSystem.class.getName());
8497
// conf.set(FileSystem.FS_DEFAULT_NAME_KEY, hdfsAddress);
8598
//https://segmentfault.com/q/1010000008473574
8699
logger.info("userHostname:{}", userHostname);
@@ -105,8 +118,9 @@ public Configuration getConfiguration() {
105118
}
106119

107120
URL endpointUrl = new URL(endpoint);
108-
conf.set("fs.s3a.endpoint", String.valueOf(endpointUrl));
109-
conf.setBoolean("fs.s3a.connection.ssl.enabled", "https".equalsIgnoreCase(endpointUrl.getProtocol()));
121+
conf.set(Constants.ENDPOINT, String.valueOf(endpointUrl));
122+
conf.setBoolean(Constants.SECURE_CONNECTIONS, "https".equalsIgnoreCase(endpointUrl.getProtocol()));
123+
conf.setBoolean(Constants.PATH_STYLE_ACCESS, pathStyleAccess);
110124

111125
userToken.accept(new IUserTokenVisitor<Void>() {
112126
@Override
@@ -116,10 +130,11 @@ public Void visit(IOffUserToken token) throws Exception {
116130

117131
@Override
118132
public Void visit(IUserNamePasswordUserToken token) throws Exception {
119-
token.getPassword();
120-
token.getPassword();
121-
conf.set("fs.s3a.access.key", token.getUserName());
122-
conf.set("fs.s3a.secret.key", token.getPassword());
133+
// token.getPassword();
134+
// token.getPassword();
135+
136+
conf.set(Constants.ACCESS_KEY, token.getUserName());
137+
conf.set(Constants.SECRET_KEY, token.getPassword());
123138
return null;
124139
// throw new UnsupportedOperationException();
125140
}
Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
1+
package com.qlangtech.tis.plugin.amazon.s3;
2+
3+
import com.google.common.collect.Lists;
4+
import com.qlangtech.tis.manage.common.Option;
5+
import com.qlangtech.tis.offline.FileSystemFactory;
6+
import org.apache.hadoop.conf.Configuration;
7+
8+
import java.util.List;
9+
import java.util.function.BiConsumer;
10+
11+
/**
12+
*
13+
* @author 百岁 (baisui@qlangtech.com)
14+
* @date 2026/3/10
15+
*/
16+
public class ReplayConfiguration extends Configuration implements FileSystemFactory.IReplayConfiguration {
17+
private final List<Option> opts = Lists.newArrayList();
18+
19+
@Override
20+
public void set(String name, String value) {
21+
super.set(name, value);
22+
opts.add(new Option(name, value));
23+
}
24+
25+
@Override
26+
public void replay(BiConsumer<String, String> consumer) {
27+
opts.forEach((opt) -> consumer.accept(opt.getName(), (String) opt.getValue()));
28+
}
29+
}
Lines changed: 28 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,30 @@
11
{
2-
2+
"endpoint": {
3+
"label": "服务端点",
4+
"placeholder": "http://localhost:9000",
5+
"help": "S3/MinIO服务端点(API地址),例如:http://localhost:9000"
6+
},
7+
"bucket": {
8+
"label": "存储桶",
9+
"help": "存储桶名称,相当于一个顶级存储容器"
10+
},
11+
"rootDir": {
12+
"label": "根目录",
13+
"placeholder": "/user/admin",
14+
"help": "系统会将数据存储到该子目录下,用户需要保证该路径有读/写权限"
15+
},
16+
"region": {
17+
"label": "区域",
18+
"placeholder": "us-east-1",
19+
"help": "可选项,如果遇到Region相关报错可以设置一个默认Region"
20+
},
21+
"userToken": {
22+
"dftVal": "userPass",
23+
"label": "认证方式",
24+
"help": "S3/MinIO认证信息,用户名填写AccessKey,密码填写SecretKey"
25+
},
26+
"pathStyleAccess": {
27+
"label": "路径访问模式",
28+
"dftVal": false
29+
}
330
}
Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
## pathStyleAccess
2+
控制 S3 API 的 URL 寻址方式。
3+
* 是(Path-Style):URL 格式为 http://endpoint/bucket/key,MinIO、Ceph 等自建 S3 兼容存储必须开启此选项。
4+
* 否(Virtual-Hosted-Style):URL 格式为 http://bucket.endpoint/key,Amazon S3 推荐使用此模式(AWS 自 2020 年 9 月起对新建 Bucket 已弃用 Path-Style)。
5+
如不确定,自建存储选 true,AWS S3 选 false。

tis-incr/tis-flink-dependency/pom.xml

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,8 @@
6060
<version>${flink.version}</version>
6161
</dependency>
6262

63+
64+
6365
<dependency>
6466
<groupId>org.apache.flink</groupId>
6567
<artifactId>flink-connector-base</artifactId>

tis-incr/tis-flink-extends/pom.xml

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,11 @@
3535

3636
<dependencies>
3737

38+
<!-- <dependency>-->
39+
<!-- <groupId>org.apache.flink</groupId>-->
40+
<!-- <artifactId>flink-s3-fs-hadoop</artifactId>-->
41+
<!-- <version>${flink.version}</version>-->
42+
<!-- </dependency>-->
3843

3944
<dependency>
4045
<groupId>org.apache.flink</groupId>

0 commit comments

Comments
 (0)