Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Contribute to GitLab
Sign in / Register
Toggle navigation
D
dlink
Project
Project
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
zhaowei
dlink
Commits
904af902
Commit
904af902
authored
Apr 20, 2022
by
wenmo
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
[Feature-435][client,executor] CDCSource sync doris
parent
4bbd40ac
Changes
13
Expand all
Hide whitespace changes
Inline
Side-by-side
Showing
13 changed files
with
996 additions
and
5 deletions
+996
-5
SinkBuilderFactory.java
...-1.11/src/main/java/com/dlink/cdc/SinkBuilderFactory.java
+2
-0
DorisSinkBuilder.java
...1/src/main/java/com/dlink/cdc/doris/DorisSinkBuilder.java
+241
-0
SinkBuilderFactory.java
...-1.12/src/main/java/com/dlink/cdc/SinkBuilderFactory.java
+2
-0
DorisSinkBuilder.java
...2/src/main/java/com/dlink/cdc/doris/DorisSinkBuilder.java
+241
-0
SinkBuilderFactory.java
...-1.13/src/main/java/com/dlink/cdc/SinkBuilderFactory.java
+2
-0
DorisSinkBuilder.java
...3/src/main/java/com/dlink/cdc/doris/DorisSinkBuilder.java
+241
-0
SinkBuilderFactory.java
...-1.14/src/main/java/com/dlink/cdc/SinkBuilderFactory.java
+2
-0
DorisSinkBuilder.java
...4/src/main/java/com/dlink/cdc/doris/DorisSinkBuilder.java
+241
-0
CDCSource.java
...executor/src/main/java/com/dlink/trans/ddl/CDCSource.java
+4
-4
pom.xml
dlink-flink/dlink-flink-1.11/pom.xml
+5
-1
pom.xml
dlink-flink/dlink-flink-1.12/pom.xml
+5
-0
pom.xml
dlink-flink/dlink-flink-1.13/pom.xml
+5
-0
pom.xml
dlink-flink/dlink-flink-1.14/pom.xml
+5
-0
No files found.
dlink-client/dlink-client-1.11/src/main/java/com/dlink/cdc/SinkBuilderFactory.java
View file @
904af902
package
com
.
dlink
.
cdc
;
package
com
.
dlink
.
cdc
;
import
com.dlink.assertion.Asserts
;
import
com.dlink.assertion.Asserts
;
import
com.dlink.cdc.doris.DorisSinkBuilder
;
import
com.dlink.cdc.kafka.KafkaSinkBuilder
;
import
com.dlink.cdc.kafka.KafkaSinkBuilder
;
import
com.dlink.exception.FlinkClientException
;
import
com.dlink.exception.FlinkClientException
;
import
com.dlink.model.FlinkCDCConfig
;
import
com.dlink.model.FlinkCDCConfig
;
...
@@ -15,6 +16,7 @@ public class SinkBuilderFactory {
...
@@ -15,6 +16,7 @@ public class SinkBuilderFactory {
private
static
SinkBuilder
[]
sinkBuilders
=
{
private
static
SinkBuilder
[]
sinkBuilders
=
{
new
KafkaSinkBuilder
(),
new
KafkaSinkBuilder
(),
new
DorisSinkBuilder
()
};
};
public
static
SinkBuilder
buildSinkBuilder
(
FlinkCDCConfig
config
)
{
public
static
SinkBuilder
buildSinkBuilder
(
FlinkCDCConfig
config
)
{
...
...
dlink-client/dlink-client-1.11/src/main/java/com/dlink/cdc/doris/DorisSinkBuilder.java
0 → 100644
View file @
904af902
This diff is collapsed.
Click to expand it.
dlink-client/dlink-client-1.12/src/main/java/com/dlink/cdc/SinkBuilderFactory.java
View file @
904af902
package
com
.
dlink
.
cdc
;
package
com
.
dlink
.
cdc
;
import
com.dlink.assertion.Asserts
;
import
com.dlink.assertion.Asserts
;
import
com.dlink.cdc.doris.DorisSinkBuilder
;
import
com.dlink.cdc.kafka.KafkaSinkBuilder
;
import
com.dlink.cdc.kafka.KafkaSinkBuilder
;
import
com.dlink.exception.FlinkClientException
;
import
com.dlink.exception.FlinkClientException
;
import
com.dlink.model.FlinkCDCConfig
;
import
com.dlink.model.FlinkCDCConfig
;
...
@@ -15,6 +16,7 @@ public class SinkBuilderFactory {
...
@@ -15,6 +16,7 @@ public class SinkBuilderFactory {
private
static
SinkBuilder
[]
sinkBuilders
=
{
private
static
SinkBuilder
[]
sinkBuilders
=
{
new
KafkaSinkBuilder
(),
new
KafkaSinkBuilder
(),
new
DorisSinkBuilder
()
};
};
public
static
SinkBuilder
buildSinkBuilder
(
FlinkCDCConfig
config
)
{
public
static
SinkBuilder
buildSinkBuilder
(
FlinkCDCConfig
config
)
{
...
...
dlink-client/dlink-client-1.12/src/main/java/com/dlink/cdc/doris/DorisSinkBuilder.java
0 → 100644
View file @
904af902
This diff is collapsed.
Click to expand it.
dlink-client/dlink-client-1.13/src/main/java/com/dlink/cdc/SinkBuilderFactory.java
View file @
904af902
package
com
.
dlink
.
cdc
;
package
com
.
dlink
.
cdc
;
import
com.dlink.assertion.Asserts
;
import
com.dlink.assertion.Asserts
;
import
com.dlink.cdc.doris.DorisSinkBuilder
;
import
com.dlink.cdc.jdbc.JdbcSinkBuilder
;
import
com.dlink.cdc.jdbc.JdbcSinkBuilder
;
import
com.dlink.cdc.kafka.KafkaSinkBuilder
;
import
com.dlink.cdc.kafka.KafkaSinkBuilder
;
import
com.dlink.exception.FlinkClientException
;
import
com.dlink.exception.FlinkClientException
;
...
@@ -17,6 +18,7 @@ public class SinkBuilderFactory {
...
@@ -17,6 +18,7 @@ public class SinkBuilderFactory {
private
static
SinkBuilder
[]
sinkBuilders
=
{
private
static
SinkBuilder
[]
sinkBuilders
=
{
new
KafkaSinkBuilder
(),
new
KafkaSinkBuilder
(),
new
JdbcSinkBuilder
(),
new
JdbcSinkBuilder
(),
new
DorisSinkBuilder
(),
};
};
public
static
SinkBuilder
buildSinkBuilder
(
FlinkCDCConfig
config
)
{
public
static
SinkBuilder
buildSinkBuilder
(
FlinkCDCConfig
config
)
{
...
...
dlink-client/dlink-client-1.13/src/main/java/com/dlink/cdc/doris/DorisSinkBuilder.java
0 → 100644
View file @
904af902
This diff is collapsed.
Click to expand it.
dlink-client/dlink-client-1.14/src/main/java/com/dlink/cdc/SinkBuilderFactory.java
View file @
904af902
package
com
.
dlink
.
cdc
;
package
com
.
dlink
.
cdc
;
import
com.dlink.assertion.Asserts
;
import
com.dlink.assertion.Asserts
;
import
com.dlink.cdc.doris.DorisSinkBuilder
;
import
com.dlink.cdc.jdbc.JdbcSinkBuilder
;
import
com.dlink.cdc.jdbc.JdbcSinkBuilder
;
import
com.dlink.cdc.kafka.KafkaSinkBuilder
;
import
com.dlink.cdc.kafka.KafkaSinkBuilder
;
import
com.dlink.exception.FlinkClientException
;
import
com.dlink.exception.FlinkClientException
;
...
@@ -17,6 +18,7 @@ public class SinkBuilderFactory {
...
@@ -17,6 +18,7 @@ public class SinkBuilderFactory {
private
static
SinkBuilder
[]
sinkBuilders
=
{
private
static
SinkBuilder
[]
sinkBuilders
=
{
new
KafkaSinkBuilder
(),
new
KafkaSinkBuilder
(),
new
JdbcSinkBuilder
(),
new
JdbcSinkBuilder
(),
new
DorisSinkBuilder
()
};
};
public
static
SinkBuilder
buildSinkBuilder
(
FlinkCDCConfig
config
)
{
public
static
SinkBuilder
buildSinkBuilder
(
FlinkCDCConfig
config
)
{
...
...
dlink-client/dlink-client-1.14/src/main/java/com/dlink/cdc/doris/DorisSinkBuilder.java
0 → 100644
View file @
904af902
This diff is collapsed.
Click to expand it.
dlink-executor/src/main/java/com/dlink/trans/ddl/CDCSource.java
View file @
904af902
...
@@ -56,9 +56,9 @@ public class CDCSource {
...
@@ -56,9 +56,9 @@ public class CDCSource {
for
(
Map
.
Entry
<
String
,
String
>
entry
:
config
.
entrySet
())
{
for
(
Map
.
Entry
<
String
,
String
>
entry
:
config
.
entrySet
())
{
if
(
entry
.
getKey
().
startsWith
(
"debezium."
))
{
if
(
entry
.
getKey
().
startsWith
(
"debezium."
))
{
String
key
=
entry
.
getKey
();
String
key
=
entry
.
getKey
();
key
=
key
.
replace
(
"debezium."
,
""
);
key
=
key
.
replace
First
(
"debezium."
,
""
);
if
(!
debezium
.
containsKey
(
key
))
{
if
(!
debezium
.
containsKey
(
key
))
{
debezium
.
put
(
entry
.
getKey
().
replace
(
"debezium."
,
""
)
,
entry
.
getValue
());
debezium
.
put
(
key
,
entry
.
getValue
());
}
}
}
}
}
}
...
@@ -66,9 +66,9 @@ public class CDCSource {
...
@@ -66,9 +66,9 @@ public class CDCSource {
for
(
Map
.
Entry
<
String
,
String
>
entry
:
config
.
entrySet
())
{
for
(
Map
.
Entry
<
String
,
String
>
entry
:
config
.
entrySet
())
{
if
(
entry
.
getKey
().
startsWith
(
"sink."
))
{
if
(
entry
.
getKey
().
startsWith
(
"sink."
))
{
String
key
=
entry
.
getKey
();
String
key
=
entry
.
getKey
();
key
=
key
.
replace
(
"sink."
,
""
);
key
=
key
.
replace
First
(
"sink."
,
""
);
if
(!
sink
.
containsKey
(
key
))
{
if
(!
sink
.
containsKey
(
key
))
{
sink
.
put
(
entry
.
getKey
().
replace
(
"sink."
,
""
)
,
entry
.
getValue
());
sink
.
put
(
key
,
entry
.
getValue
());
}
}
}
}
}
}
...
...
dlink-flink/dlink-flink-1.11/pom.xml
View file @
904af902
...
@@ -90,6 +90,10 @@
...
@@ -90,6 +90,10 @@
<groupId>
org.slf4j
</groupId>
<groupId>
org.slf4j
</groupId>
<artifactId>
slf4j-api
</artifactId>
<artifactId>
slf4j-api
</artifactId>
</dependency>
</dependency>
<dependency>
<groupId>
org.apache.doris
</groupId>
<artifactId>
flink-doris-connector-1.11_2.12
</artifactId>
<version>
1.0.3
</version>
</dependency>
</dependencies>
</dependencies>
</project>
</project>
\ No newline at end of file
dlink-flink/dlink-flink-1.12/pom.xml
View file @
904af902
...
@@ -90,6 +90,11 @@
...
@@ -90,6 +90,11 @@
<groupId>
org.slf4j
</groupId>
<groupId>
org.slf4j
</groupId>
<artifactId>
slf4j-api
</artifactId>
<artifactId>
slf4j-api
</artifactId>
</dependency>
</dependency>
<dependency>
<groupId>
org.apache.doris
</groupId>
<artifactId>
flink-doris-connector-1.12_2.12
</artifactId>
<version>
1.0.3
</version>
</dependency>
</dependencies>
</dependencies>
</project>
</project>
\ No newline at end of file
dlink-flink/dlink-flink-1.13/pom.xml
View file @
904af902
...
@@ -115,5 +115,10 @@
...
@@ -115,5 +115,10 @@
<groupId>
org.slf4j
</groupId>
<groupId>
org.slf4j
</groupId>
<artifactId>
slf4j-api
</artifactId>
<artifactId>
slf4j-api
</artifactId>
</dependency>
</dependency>
<dependency>
<groupId>
org.apache.doris
</groupId>
<artifactId>
flink-doris-connector-1.13_2.12
</artifactId>
<version>
1.0.3
</version>
</dependency>
</dependencies>
</dependencies>
</project>
</project>
\ No newline at end of file
dlink-flink/dlink-flink-1.14/pom.xml
View file @
904af902
...
@@ -101,5 +101,10 @@
...
@@ -101,5 +101,10 @@
<artifactId>
commons-cli
</artifactId>
<artifactId>
commons-cli
</artifactId>
<version>
${commons.version}
</version>
<version>
${commons.version}
</version>
</dependency>
</dependency>
<dependency>
<groupId>
org.apache.doris
</groupId>
<artifactId>
flink-doris-connector-1.14_2.12
</artifactId>
<version>
1.0.3
</version>
</dependency>
</dependencies>
</dependencies>
</project>
</project>
\ No newline at end of file
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment