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
89cfb186
Commit
89cfb186
authored
Apr 23, 2022
by
wenmo
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
[Fix-442][client,executor] Add SOURCE_AVRO_SCHEMA
parent
3e62e609
Changes
1
Hide whitespace changes
Inline
Side-by-side
Showing
1 changed file
with
5 additions
and
2 deletions
+5
-2
HudiSinkBuilder.java
....13/src/main/java/com/dlink/cdc/hudi/HudiSinkBuilder.java
+5
-2
No files found.
dlink-client/dlink-client-1.13/src/main/java/com/dlink/cdc/hudi/HudiSinkBuilder.java
View file @
89cfb186
...
@@ -10,6 +10,7 @@ import org.apache.hudi.common.model.HoodieRecord;
...
@@ -10,6 +10,7 @@ import org.apache.hudi.common.model.HoodieRecord;
import
org.apache.hudi.common.model.HoodieTableType
;
import
org.apache.hudi.common.model.HoodieTableType
;
import
org.apache.hudi.configuration.FlinkOptions
;
import
org.apache.hudi.configuration.FlinkOptions
;
import
org.apache.hudi.sink.utils.Pipelines
;
import
org.apache.hudi.sink.utils.Pipelines
;
import
org.apache.hudi.util.AvroSchemaConverter
;
import
java.io.Serializable
;
import
java.io.Serializable
;
import
java.util.List
;
import
java.util.List
;
...
@@ -17,7 +18,6 @@ import java.util.Map;
...
@@ -17,7 +18,6 @@ import java.util.Map;
import
com.dlink.cdc.AbstractSinkBuilder
;
import
com.dlink.cdc.AbstractSinkBuilder
;
import
com.dlink.cdc.SinkBuilder
;
import
com.dlink.cdc.SinkBuilder
;
import
com.dlink.cdc.doris.DorisSinkBuilder
;
import
com.dlink.model.FlinkCDCConfig
;
import
com.dlink.model.FlinkCDCConfig
;
import
com.dlink.model.Table
;
import
com.dlink.model.Table
;
...
@@ -76,8 +76,11 @@ public class HudiSinkBuilder extends AbstractSinkBuilder implements SinkBuilder,
...
@@ -76,8 +76,11 @@ public class HudiSinkBuilder extends AbstractSinkBuilder implements SinkBuilder,
configuration
.
set
(
FlinkOptions
.
TABLE_NAME
,
table
.
getSchemaTableNameWithUnderline
());
configuration
.
set
(
FlinkOptions
.
TABLE_NAME
,
table
.
getSchemaTableNameWithUnderline
());
configuration
.
set
(
FlinkOptions
.
HIVE_SYNC_DB
,
table
.
getSchema
());
configuration
.
set
(
FlinkOptions
.
HIVE_SYNC_DB
,
table
.
getSchema
());
configuration
.
set
(
FlinkOptions
.
HIVE_SYNC_TABLE
,
table
.
getName
());
configuration
.
set
(
FlinkOptions
.
HIVE_SYNC_TABLE
,
table
.
getName
());
RowType
rowType
=
RowType
.
of
(
false
,
columnTypes
,
columnNames
);
configuration
.
setString
(
FlinkOptions
.
SOURCE_AVRO_SCHEMA
,
AvroSchemaConverter
.
convertToSchema
(
rowType
,
table
.
getSchemaTableNameWithUnderline
()).
toString
());
DataStream
<
HoodieRecord
>
hoodieRecordDataStream
=
Pipelines
.
bootstrap
(
configuration
,
RowType
.
of
(
columnTypes
,
columnNames
)
,
parallelism
,
rowDataDataStream
);
DataStream
<
HoodieRecord
>
hoodieRecordDataStream
=
Pipelines
.
bootstrap
(
configuration
,
rowType
,
parallelism
,
rowDataDataStream
);
DataStream
<
Object
>
pipeline
=
Pipelines
.
hoodieStreamWrite
(
configuration
,
parallelism
,
hoodieRecordDataStream
);
DataStream
<
Object
>
pipeline
=
Pipelines
.
hoodieStreamWrite
(
configuration
,
parallelism
,
hoodieRecordDataStream
);
if
(
isMor
)
{
if
(
isMor
)
{
...
...
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