Commit e3e58132 by chenyuanjie

SparkToDoris 部分列更新公共方法

parent 6d3d2616
...@@ -236,11 +236,12 @@ class DorisHelper(object): ...@@ -236,11 +236,12 @@ class DorisHelper(object):
f"执行导出数据到doris任务,导出数据: {df_save}, doris库名: {db_name}, doris表名: {table_name}, 更新字段为: {update_field}, 更新模式为: {update_mode}") f"执行导出数据到doris任务,导出数据: {df_save}, doris库名: {db_name}, doris表名: {table_name}, 更新字段为: {update_field}, 更新模式为: {update_mode}")
try: try:
connection_info = DorisHelper.get_connection_info(use_type) connection_info = DorisHelper.get_connection_info(use_type)
df_save.write.format("starrocks") \ df_save.write.format("doris") \
.option("doris.fenodes", f"{connection_info['ip']}:{connection_info['http_port']}") \ .option("doris.fenodes", f"{connection_info['ip']}:{connection_info['http_port']}") \
.option("doris.table.identifier", f"{db_name}.{table_name}") \ .option("doris.table.identifier", f"{db_name}.{table_name}") \
.option("user", connection_info['user']) \ .option("user", connection_info['user']) \
.option("password", connection_info['pwd']) \ .option("password", connection_info['pwd']) \
.option("doris.write.fields", f"{update_field}") \
.option("doris.sink.properties.partial_columns", "true") \ .option("doris.sink.properties.partial_columns", "true") \
.option("doris.sink.batch.interval.ms", "30000") \ .option("doris.sink.batch.interval.ms", "30000") \
.option("doris.sink.properties.format", "json") \ .option("doris.sink.properties.format", "json") \
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment