forked from apache/doris
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathDataSinks.thrift
More file actions
467 lines (407 loc) · 14.1 KB
/
Copy pathDataSinks.thrift
File metadata and controls
467 lines (407 loc) · 14.1 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
namespace cpp doris
namespace java org.apache.doris.thrift
include "Exprs.thrift"
include "Types.thrift"
include "Descriptors.thrift"
include "Partitions.thrift"
include "PlanNodes.thrift"
enum TDataSinkType {
DATA_STREAM_SINK = 0,
RESULT_SINK = 1,
DATA_SPLIT_SINK = 2, // deprecated
MYSQL_TABLE_SINK = 3,
EXPORT_SINK = 4,
OLAP_TABLE_SINK = 5,
MEMORY_SCRATCH_SINK = 6,
ODBC_TABLE_SINK = 7,
RESULT_FILE_SINK = 8,
JDBC_TABLE_SINK = 9,
MULTI_CAST_DATA_STREAM_SINK = 10,
GROUP_COMMIT_OLAP_TABLE_SINK = 11, // deprecated
GROUP_COMMIT_BLOCK_SINK = 12,
HIVE_TABLE_SINK = 13,
ICEBERG_TABLE_SINK = 14,
DICTIONARY_SINK = 15,
BLACKHOLE_SINK = 16,
}
enum TResultSinkType {
MYSQL_PROTOCOL = 0,
ARROW_FLIGHT_PROTOCOL = 1,
FILE = 2, // deprecated, should not be used any more. FileResultSink is covered by TRESULT_FILE_SINK for concurrent purpose.
}
enum TParquetCompressionType {
SNAPPY = 0,
GZIP = 1,
BROTLI = 2,
ZSTD = 3,
LZ4 = 4,
LZO = 5,
BZ2 = 6,
UNCOMPRESSED = 7,
}
enum TParquetVersion {
PARQUET_1_0 = 0,
PARQUET_2_LATEST = 1,
}
enum TParquetDataType {
BOOLEAN = 0,
INT32 = 1,
INT64 = 2,
INT96 = 3,
BYTE_ARRAY = 4,
FLOAT = 5,
DOUBLE = 6,
FIXED_LEN_BYTE_ARRAY = 7,
}
enum TParquetDataLogicalType {
UNDEFINED = 0, // Not a real logical type
STRING = 1,
MAP = 2,
LIST = 3,
ENUM = 4,
DECIMAL = 5,
DATE = 6,
TIME = 7,
TIMESTAMP = 8,
INTERVAL = 9,
INT = 10,
NIL = 11, // Thrift NullType: annotates data that is always null
JSON = 12,
BSON = 13,
UUID = 14,
NONE = 15 // Not a real logical type; should always be last element
}
enum TParquetRepetitionType {
REQUIRED = 0,
REPEATED = 1,
OPTIONAL = 2,
}
struct TParquetSchema {
1: optional TParquetRepetitionType schema_repetition_type
2: optional TParquetDataType schema_data_type
3: optional string schema_column_name
4: optional TParquetDataLogicalType schema_data_logical_type
}
struct TResultFileSinkOptions {
1: required string file_path
2: required PlanNodes.TFileFormatType file_format
3: optional string column_separator // only for csv
4: optional string line_delimiter // only for csv
5: optional i64 max_file_size_bytes
6: optional list<Types.TNetworkAddress> broker_addresses; // only for remote file
7: optional map<string, string> broker_properties // only for remote file
8: optional string success_file_name
9: optional list<list<string>> schema // for orc file
10: optional map<string, string> file_properties // for orc file
//note: use outfile with parquet format, have deprecated 9:schema and 10:file_properties
//because when this info thrift to BE, BE hava to find useful info in string,
//have to check by use string directly, and maybe not so efficient
11: optional list<TParquetSchema> parquet_schemas
12: optional TParquetCompressionType parquet_compression_type
13: optional bool parquet_disable_dictionary
14: optional TParquetVersion parquet_version
15: optional string orc_schema
16: optional bool delete_existing_files;
17: optional string file_suffix;
18: optional bool with_bom;
19: optional PlanNodes.TFileCompressType orc_compression_type;
// Since we have changed the type mapping from Doris to Orc type,
// using the Outfile to export Date/Datetime types will cause BE core dump
// when only upgrading BE without upgrading FE.
// orc_writer_version = 1 means doris FE is higher than version 2.1.5
// orc_writer_version = 0 means doris FE is less than or equal to version 2.1.5
20: optional i64 orc_writer_version;
//iceberg write sink use int64
//hive write sink use int96
//export data to file use by user define properties
21: optional bool enable_int96_timestamps
// currently only for csv
// TODO: merge with parquet_compression_type and orc_compression_type
22: optional PlanNodes.TFileCompressType compression_type
}
struct TMemoryScratchSink {
}
// Specification of one output destination of a plan fragment
struct TPlanFragmentDestination {
// the globally unique fragment instance id
1: required Types.TUniqueId fragment_instance_id
// ... which is being executed on this server
2: required Types.TNetworkAddress server
3: optional Types.TNetworkAddress brpc_server
}
// Sink which forwards data to a remote plan fragment,
// according to the given output partition specification
// (ie, the m:1 part of an m:n data stream)
struct TDataStreamSink {
// destination node id
1: required Types.TPlanNodeId dest_node_id
// Specification of how the output of a fragment is partitioned.
// If the partitioning type is UNPARTITIONED, the output is broadcast
// to each destination host.
2: required Partitions.TDataPartition output_partition
3: optional bool ignore_not_found
// per-destination projections
4: optional list<Exprs.TExpr> output_exprs
// project output tuple id
5: optional Types.TTupleId output_tuple_id
// per-destination filters
6: optional list<Exprs.TExpr> conjuncts
// per-destination runtime filters
7: optional list<PlanNodes.TRuntimeFilterDesc> runtime_filters
// used for partition_type = OLAP_TABLE_SINK_HASH_PARTITIONED
8: optional Descriptors.TOlapTableSchemaParam tablet_sink_schema
9: optional Descriptors.TOlapTablePartitionParam tablet_sink_partition
10: optional Descriptors.TOlapTableLocationParam tablet_sink_location
11: optional i64 tablet_sink_txn_id
12: optional Types.TTupleId tablet_sink_tuple_id
13: optional list<Exprs.TExpr> tablet_sink_exprs
14: optional bool is_merge
}
struct TMultiCastDataStreamSink {
1: optional list<TDataStreamSink> sinks;
2: optional list<list<TPlanFragmentDestination>> destinations;
}
struct TFetchOption {
1: optional bool use_two_phase_fetch;
// Nodes in this cluster, used for second phase fetch
2: optional Descriptors.TPaloNodesInfo nodes_info;
// Whether fetch row store
3: optional bool fetch_row_store;
// Fetch schema
4: optional list<Descriptors.TColumn> column_desc;
}
struct TResultSink {
1: optional TResultSinkType type;
2: optional TResultFileSinkOptions file_options; // deprecated
3: optional TFetchOption fetch_option;
}
struct TResultFileSink {
1: optional TResultFileSinkOptions file_options;
2: optional Types.TStorageBackendType storage_backend_type;
3: optional Types.TPlanNodeId dest_node_id;
4: optional Types.TTupleId output_tuple_id;
5: optional string header;
6: optional string header_type;
}
struct TMysqlTableSink {
1: required string host
2: required i32 port
3: required string user
4: required string passwd
5: required string db
6: required string table
7: required string charset
}
struct TOdbcTableSink {
1: optional string connect_string
2: optional string table
3: optional bool use_transaction
}
struct TJdbcTableSink {
1: optional Descriptors.TJdbcTable jdbc_table
2: optional bool use_transaction
3: optional Types.TOdbcTableType table_type
4: optional string insert_sql
}
struct TExportSink {
1: required Types.TFileType file_type
2: required string export_path
3: required string column_separator
4: required string line_delimiter
// properties need to access broker.
5: optional list<Types.TNetworkAddress> broker_addresses
6: optional map<string, string> properties
7: optional string header
}
enum TGroupCommitMode {
SYNC_MODE = 0,
ASYNC_MODE = 1,
OFF_MODE = 2
}
struct TOlapTableSink {
1: required Types.TUniqueId load_id
2: required i64 txn_id
3: required i64 db_id
4: required i64 table_id
5: required i32 tuple_id
6: required i32 num_replicas
7: required bool need_gen_rollup // Deprecated, not used since alter job v2
8: optional string db_name
9: optional string table_name
10: required Descriptors.TOlapTableSchemaParam schema
11: required Descriptors.TOlapTablePartitionParam partition
12: required Descriptors.TOlapTableLocationParam location
13: required Descriptors.TPaloNodesInfo nodes_info
14: optional i64 load_channel_timeout_s // the timeout of load channels in second
15: optional i32 send_batch_parallelism
16: optional bool load_to_single_tablet
17: optional bool write_single_replica
18: optional Descriptors.TOlapTableLocationParam slave_location
19: optional i64 txn_timeout_s // timeout of load txn in second
20: optional bool write_file_cache
// used by GroupCommitBlockSink
21: optional i64 base_schema_version
22: optional TGroupCommitMode group_commit_mode
23: optional double max_filter_ratio
24: optional string storage_vault_id
}
struct THiveLocationParams {
1: optional string write_path
2: optional string target_path
3: optional Types.TFileType file_type
// Other object store will convert write_path to s3 scheme path for BE, this field keeps the original write path.
4: optional string original_write_path
}
struct TSortedColumn {
1: optional string sort_column_name
2: optional i32 order // asc(1) or desc(0)
}
struct TBucketingMode {
1: optional i32 bucket_version
}
struct THiveBucket {
1: optional list<string> bucketed_by
2: optional TBucketingMode bucket_mode
3: optional i32 bucket_count
4: optional list<TSortedColumn> sorted_by
}
enum THiveColumnType {
PARTITION_KEY = 0,
REGULAR = 1,
SYNTHESIZED = 2
}
struct THiveColumn {
1: optional string name
2: optional THiveColumnType column_type
}
struct THivePartition {
1: optional list<string> values
2: optional THiveLocationParams location
3: optional PlanNodes.TFileFormatType file_format
}
struct THiveSerDeProperties {
1: optional string field_delim
2: optional string line_delim
3: optional string collection_delim // array ,map ,struct delimiter
4: optional string mapkv_delim
5: optional string escape_char
6: optional string null_format
}
struct THiveTableSink {
1: optional string db_name
2: optional string table_name
3: optional list<THiveColumn> columns
4: optional list<THivePartition> partitions
5: optional THiveBucket bucket_info
6: optional PlanNodes.TFileFormatType file_format
7: optional PlanNodes.TFileCompressType compression_type
8: optional THiveLocationParams location
9: optional map<string, string> hadoop_config
10: optional bool overwrite
11: optional THiveSerDeProperties serde_properties
12: optional list<Types.TNetworkAddress> broker_addresses;
}
enum TUpdateMode {
NEW = 0, // add partition
APPEND = 1, // alter partition
OVERWRITE = 2 // insert overwrite
}
struct TS3MPUPendingUpload {
1: optional string bucket
2: optional string key
3: optional string upload_id
4: optional map<i32, string> etags
}
struct THivePartitionUpdate {
1: optional string name
2: optional TUpdateMode update_mode
3: optional THiveLocationParams location
4: optional list<string> file_names
5: optional i64 row_count
6: optional i64 file_size
7: optional list<TS3MPUPendingUpload> s3_mpu_pending_uploads
}
enum TFileContent {
DATA = 0,
POSITION_DELETES = 1,
EQUALITY_DELETES = 2
}
struct TIcebergCommitData {
1: optional string file_path
2: optional i64 row_count
3: optional i64 file_size
4: optional TFileContent file_content
5: optional list<string> partition_values
6: optional list<string> referenced_data_files
}
struct TSortField {
1: optional i32 source_column_id
2: optional bool ascending
3: optional bool null_first
}
struct TIcebergTableSink {
1: optional string db_name
2: optional string tb_name
3: optional string schema_json
4: optional map<i32, string> partition_specs_json
5: optional i32 partition_spec_id
6: optional list<TSortField> sort_fields
7: optional PlanNodes.TFileFormatType file_format
8: optional string output_path
9: optional map<string, string> hadoop_config
10: optional bool overwrite
11: optional Types.TFileType file_type
12: optional string original_output_path
13: optional PlanNodes.TFileCompressType compression_type
14: optional list<Types.TNetworkAddress> broker_addresses;
}
enum TDictLayoutType {
HASH_MAP = 0,
IP_TRIE = 1,
}
struct TDictionarySink {
1: optional i64 dictionary_id
2: optional i64 version_id
3: optional string dictionary_name
4: optional TDictLayoutType layout_type
5: optional list<i64> key_output_expr_slots
6: optional list<i64> value_output_expr_slots
7: optional list<string> value_names
8: optional bool skip_null_key
9: optional i64 memory_limit
}
struct TBlackholeSink {
}
struct TDataSink {
1: required TDataSinkType type
2: optional TDataStreamSink stream_sink
3: optional TResultSink result_sink
5: optional TMysqlTableSink mysql_table_sink
6: optional TExportSink export_sink
7: optional TOlapTableSink olap_table_sink
8: optional TMemoryScratchSink memory_scratch_sink
9: optional TOdbcTableSink odbc_table_sink
10: optional TResultFileSink result_file_sink
11: optional TJdbcTableSink jdbc_table_sink
12: optional TMultiCastDataStreamSink multi_cast_stream_sink
13: optional THiveTableSink hive_table_sink
14: optional TIcebergTableSink iceberg_table_sink
15: optional TDictionarySink dictionary_sink
16: optional TBlackholeSink blackhole_sink
}