Coverage Report

Created: 2026-07-21 18:02

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exec/sink/viceberg_merge_sink.cpp
Line
Count
Source
1
// Licensed to the Apache Software Foundation (ASF) under one
2
// or more contributor license agreements.  See the NOTICE file
3
// distributed with this work for additional information
4
// regarding copyright ownership.  The ASF licenses this file
5
// to you under the Apache License, Version 2.0 (the
6
// "License"); you may not use this file except in compliance
7
// with the License.  You may obtain a copy of the License at
8
//
9
//   http://www.apache.org/licenses/LICENSE-2.0
10
//
11
// Unless required by applicable law or agreed to in writing,
12
// software distributed under the License is distributed on an
13
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14
// KIND, either express or implied.  See the License for the
15
// specific language governing permissions and limitations
16
// under the License.
17
18
#include "exec/sink/viceberg_merge_sink.h"
19
20
#include <fmt/format.h>
21
22
#include "common/consts.h"
23
#include "common/exception.h"
24
#include "common/logging.h"
25
#include "core/block/block.h"
26
#include "core/column/column_nullable.h"
27
#include "core/column/column_vector.h"
28
#include "exec/sink/sink_common.h"
29
#include "exec/sink/viceberg_delete_sink.h"
30
#include "exec/sink/writer/iceberg/viceberg_table_writer.h"
31
#include "exprs/vexpr_context.h"
32
#include "format/table/iceberg/schema.h"
33
#include "format/table/iceberg/schema_parser.h"
34
#include "runtime/runtime_state.h"
35
#include "util/string_util.h"
36
37
namespace doris {
38
39
namespace {} // namespace
40
41
VIcebergMergeSink::VIcebergMergeSink(const TDataSink& t_sink, const VExprContextSPtrs& output_exprs,
42
                                     std::shared_ptr<Dependency> dep,
43
                                     std::shared_ptr<Dependency> fin_dep)
44
1.51k
        : AsyncResultWriter(output_exprs, dep, fin_dep), _t_sink(t_sink) {
45
1.51k
    DCHECK(_t_sink.__isset.iceberg_merge_sink);
46
1.51k
}
47
48
1.51k
VIcebergMergeSink::~VIcebergMergeSink() = default;
49
50
1.51k
Status VIcebergMergeSink::init_properties(ObjectPool* pool, const RowDescriptor& row_desc) {
51
1.51k
    RETURN_IF_ERROR(_build_inner_sinks());
52
53
1.51k
    _table_writer = std::make_unique<VIcebergTableWriter>(_table_sink, _table_output_expr_ctxs,
54
1.51k
                                                          nullptr, nullptr);
55
1.51k
    _delete_writer = std::make_unique<VIcebergDeleteSink>(_delete_sink, _delete_output_expr_ctxs,
56
1.51k
                                                          nullptr, nullptr);
57
1.51k
    RETURN_IF_ERROR(_table_writer->init_properties(pool, row_desc));
58
1.51k
    RETURN_IF_ERROR(_delete_writer->init_properties(pool));
59
1.51k
    return Status::OK();
60
1.51k
}
61
62
1.51k
Status VIcebergMergeSink::open(RuntimeState* state, RuntimeProfile* profile) {
63
1.51k
    _state = state;
64
65
1.51k
    _written_rows_counter = ADD_COUNTER(profile, "RowsWritten", TUnit::UNIT);
66
1.51k
    _insert_rows_counter = ADD_COUNTER(profile, "InsertRows", TUnit::UNIT);
67
1.51k
    _delete_rows_counter = ADD_COUNTER(profile, "DeleteRows", TUnit::UNIT);
68
1.51k
    _send_data_timer = ADD_TIMER(profile, "SendDataTime");
69
1.51k
    _open_timer = ADD_TIMER(profile, "OpenTime");
70
1.51k
    _close_timer = ADD_TIMER(profile, "CloseTime");
71
72
1.51k
    SCOPED_TIMER(_open_timer);
73
74
1.51k
    RETURN_IF_ERROR(_prepare_output_layout());
75
76
1.50k
    RuntimeProfile* table_profile = profile->create_child("IcebergMergeTableWriter", true, true);
77
1.50k
    RuntimeProfile* delete_profile = profile->create_child("IcebergMergeDeleteWriter", true, true);
78
79
1.50k
    RETURN_IF_ERROR(_table_writer->open(state, table_profile));
80
1.50k
    RETURN_IF_ERROR(_delete_writer->open(state, delete_profile));
81
82
1.50k
    return Status::OK();
83
1.50k
}
84
85
377
Status VIcebergMergeSink::write(RuntimeState* state, Block& block) {
86
377
    SCOPED_TIMER(_send_data_timer);
87
377
    if (block.rows() == 0) {
88
0
        return Status::OK();
89
0
    }
90
91
377
    Block output_block;
92
377
    RETURN_IF_ERROR(_projection_block(block, &output_block));
93
377
    if (output_block.rows() == 0) {
94
0
        return Status::OK();
95
0
    }
96
97
377
    _row_count += output_block.rows();
98
99
377
    if (_operation_idx < 0 || _row_id_idx < 0) {
100
0
        return Status::InternalError("Iceberg merge sink missing operation/row_id columns");
101
0
    }
102
103
377
    const auto& op_column = output_block.get_by_position(_operation_idx).column;
104
377
    const auto* op_data = remove_nullable(op_column).get();
105
106
377
    IColumn::Filter delete_filter(output_block.rows(), 0);
107
377
    IColumn::Filter insert_filter(output_block.rows(), 0);
108
377
    bool has_delete = false;
109
377
    bool has_insert = false;
110
377
    size_t delete_rows = 0;
111
377
    size_t insert_rows = 0;
112
113
854
    for (size_t i = 0; i < output_block.rows(); ++i) {
114
478
        int8_t op = static_cast<int8_t>(op_data->get_int(i));
115
478
        bool delete_op = is_delete_op(op);
116
478
        bool insert_op = is_insert_op(op);
117
478
        if (!delete_op && !insert_op) {
118
1
            return Status::InternalError("Unknown Iceberg merge operation {}", op);
119
1
        }
120
477
        if (delete_op) {
121
239
            delete_filter[i] = 1;
122
239
            has_delete = true;
123
239
            ++_delete_row_count;
124
239
            ++delete_rows;
125
239
        }
126
477
        if (insert_op) {
127
239
            insert_filter[i] = 1;
128
239
            has_insert = true;
129
239
            ++_insert_row_count;
130
239
            ++insert_rows;
131
239
        }
132
477
    }
133
134
376
    bool skip_io = false;
135
#ifdef BE_TEST
136
    skip_io = _skip_io;
137
#endif
138
139
376
    if (has_delete && !skip_io) {
140
216
        Block delete_block = output_block;
141
216
        std::vector<int> delete_indices {_row_id_idx};
142
216
        delete_block.erase_not_in(delete_indices);
143
216
        Block::filter_block_internal(&delete_block, delete_filter);
144
216
        RETURN_IF_ERROR(_delete_writer->write(state, delete_block));
145
216
    }
146
147
376
    if (has_insert && !skip_io) {
148
196
        if (_data_column_indices.empty()) {
149
0
            return Status::InternalError("Iceberg merge sink has no data columns for insert");
150
0
        }
151
196
        Block insert_block = output_block;
152
196
        insert_block.erase_not_in(_data_column_indices);
153
196
        Block::filter_block_internal(&insert_block, insert_filter);
154
196
        RETURN_IF_ERROR(_table_writer->write_prepared_block(insert_block));
155
196
    }
156
157
376
    if (_written_rows_counter != nullptr) {
158
376
        COUNTER_UPDATE(_written_rows_counter, output_block.rows());
159
376
    }
160
376
    if (_insert_rows_counter != nullptr) {
161
376
        COUNTER_UPDATE(_insert_rows_counter, insert_rows);
162
376
    }
163
376
    if (_delete_rows_counter != nullptr) {
164
376
        COUNTER_UPDATE(_delete_rows_counter, delete_rows);
165
376
    }
166
167
376
    return Status::OK();
168
376
}
169
170
1.50k
Status VIcebergMergeSink::close(Status close_status) {
171
1.50k
    SCOPED_TIMER(_close_timer);
172
173
1.50k
    if (!close_status.ok()) {
174
0
        LOG(WARNING) << fmt::format("VIcebergMergeSink close with error: {}",
175
0
                                    close_status.to_string());
176
0
        if (_table_writer) {
177
0
            static_cast<void>(_table_writer->close(close_status));
178
0
        }
179
0
        if (_delete_writer) {
180
0
            static_cast<void>(_delete_writer->close(close_status));
181
0
        }
182
0
        return close_status;
183
0
    }
184
185
1.50k
    Status table_status = Status::OK();
186
1.50k
    Status delete_status = Status::OK();
187
1.50k
    if (_table_writer) {
188
1.50k
        table_status = _table_writer->close(close_status);
189
1.50k
    }
190
1.50k
    if (_delete_writer) {
191
1.50k
        delete_status = _delete_writer->close(close_status);
192
1.50k
    }
193
194
1.50k
    if (_written_rows_counter != nullptr) {
195
1.50k
        COUNTER_SET(_written_rows_counter, static_cast<int64_t>(_row_count));
196
1.50k
    }
197
1.50k
    if (_insert_rows_counter != nullptr) {
198
1.50k
        COUNTER_SET(_insert_rows_counter, static_cast<int64_t>(_insert_row_count));
199
1.50k
    }
200
1.50k
    if (_delete_rows_counter != nullptr) {
201
1.50k
        COUNTER_SET(_delete_rows_counter, static_cast<int64_t>(_delete_row_count));
202
1.50k
    }
203
204
1.50k
    if (!table_status.ok()) {
205
0
        return table_status;
206
0
    }
207
1.50k
    return delete_status;
208
1.50k
}
209
210
1.51k
Status VIcebergMergeSink::_build_inner_sinks() {
211
1.51k
    if (!_t_sink.__isset.iceberg_merge_sink) {
212
0
        return Status::InternalError("Missing iceberg merge sink config");
213
0
    }
214
215
1.51k
    const auto& merge_sink = _t_sink.iceberg_merge_sink;
216
217
1.51k
    TIcebergTableSink table_sink;
218
1.51k
    if (merge_sink.__isset.db_name) {
219
1.51k
        table_sink.__set_db_name(merge_sink.db_name);
220
1.51k
    }
221
1.51k
    if (merge_sink.__isset.tb_name) {
222
1.51k
        table_sink.__set_tb_name(merge_sink.tb_name);
223
1.51k
    }
224
1.51k
    if (merge_sink.__isset.schema_json) {
225
1.51k
        table_sink.__set_schema_json(merge_sink.schema_json);
226
1.51k
    }
227
1.51k
    if (merge_sink.__isset.partition_specs_json) {
228
608
        table_sink.__set_partition_specs_json(merge_sink.partition_specs_json);
229
608
    }
230
1.51k
    if (merge_sink.__isset.partition_spec_id) {
231
614
        table_sink.__set_partition_spec_id(merge_sink.partition_spec_id);
232
614
    }
233
1.51k
    if (merge_sink.__isset.sort_fields) {
234
0
        table_sink.__set_sort_fields(merge_sink.sort_fields);
235
0
    }
236
1.51k
    if (merge_sink.__isset.file_format) {
237
1.51k
        table_sink.__set_file_format(merge_sink.file_format);
238
1.51k
    }
239
1.51k
    if (merge_sink.__isset.compression_type) {
240
1.51k
        table_sink.__set_compression_type(merge_sink.compression_type);
241
1.51k
    }
242
1.51k
    if (merge_sink.__isset.output_path) {
243
1.51k
        table_sink.__set_output_path(merge_sink.output_path);
244
1.51k
    }
245
1.51k
    if (merge_sink.__isset.original_output_path) {
246
1.51k
        table_sink.__set_original_output_path(merge_sink.original_output_path);
247
1.51k
    }
248
1.51k
    if (merge_sink.__isset.hadoop_config) {
249
1.50k
        table_sink.__set_hadoop_config(merge_sink.hadoop_config);
250
1.50k
    }
251
1.51k
    if (merge_sink.__isset.file_type) {
252
1.51k
        table_sink.__set_file_type(merge_sink.file_type);
253
1.51k
    }
254
1.51k
    if (merge_sink.__isset.broker_addresses) {
255
0
        table_sink.__set_broker_addresses(merge_sink.broker_addresses);
256
0
    }
257
1.51k
    if (merge_sink.__isset.collect_column_stats) {
258
1.50k
        table_sink.__set_collect_column_stats(merge_sink.collect_column_stats);
259
1.50k
    }
260
1.51k
    _table_sink.__set_type(TDataSinkType::ICEBERG_TABLE_SINK);
261
1.51k
    _table_sink.__set_iceberg_table_sink(table_sink);
262
263
1.51k
    TIcebergDeleteSink delete_sink;
264
1.51k
    if (merge_sink.__isset.db_name) {
265
1.51k
        delete_sink.__set_db_name(merge_sink.db_name);
266
1.51k
    }
267
1.51k
    if (merge_sink.__isset.tb_name) {
268
1.51k
        delete_sink.__set_tb_name(merge_sink.tb_name);
269
1.51k
    }
270
1.51k
    if (merge_sink.__isset.delete_type) {
271
1.51k
        delete_sink.__set_delete_type(merge_sink.delete_type);
272
1.51k
    }
273
1.51k
    if (merge_sink.__isset.file_format) {
274
1.51k
        delete_sink.__set_file_format(merge_sink.file_format);
275
1.51k
    }
276
1.51k
    if (merge_sink.__isset.compression_type) {
277
1.51k
        delete_sink.__set_compress_type(merge_sink.compression_type);
278
1.51k
    }
279
1.51k
    if (merge_sink.__isset.output_path) {
280
1.51k
        delete_sink.__set_output_path(merge_sink.output_path);
281
1.51k
    }
282
1.51k
    if (merge_sink.__isset.table_location) {
283
1.51k
        delete_sink.__set_table_location(merge_sink.table_location);
284
1.51k
    }
285
1.51k
    if (merge_sink.__isset.hadoop_config) {
286
1.50k
        delete_sink.__set_hadoop_config(merge_sink.hadoop_config);
287
1.50k
    }
288
1.51k
    if (merge_sink.__isset.file_type) {
289
1.51k
        delete_sink.__set_file_type(merge_sink.file_type);
290
1.51k
    }
291
1.51k
    if (merge_sink.__isset.partition_spec_id_for_delete) {
292
614
        delete_sink.__set_partition_spec_id(merge_sink.partition_spec_id_for_delete);
293
614
    }
294
1.51k
    if (merge_sink.__isset.partition_data_json_for_delete) {
295
0
        delete_sink.__set_partition_data_json(merge_sink.partition_data_json_for_delete);
296
0
    }
297
1.51k
    if (merge_sink.__isset.broker_addresses) {
298
0
        delete_sink.__set_broker_addresses(merge_sink.broker_addresses);
299
0
    }
300
1.51k
    if (merge_sink.__isset.format_version) {
301
1.50k
        delete_sink.__set_format_version(merge_sink.format_version);
302
1.50k
    }
303
1.51k
    if (merge_sink.__isset.rewritable_delete_file_sets) {
304
1.50k
        delete_sink.__set_rewritable_delete_file_sets(merge_sink.rewritable_delete_file_sets);
305
1.50k
    }
306
1.51k
    _delete_sink.__set_type(TDataSinkType::ICEBERG_DELETE_SINK);
307
1.51k
    _delete_sink.__set_iceberg_delete_sink(delete_sink);
308
309
1.51k
    return Status::OK();
310
1.51k
}
311
312
1.50k
Status VIcebergMergeSink::_prepare_output_layout() {
313
1.50k
    if (_vec_output_expr_ctxs.empty()) {
314
0
        return Status::InternalError("Iceberg merge sink has empty output expressions");
315
0
    }
316
317
1.50k
    std::string row_id_name = doris::to_lower(BeConsts::ICEBERG_ROWID_COL);
318
1.50k
    std::string op_name = doris::to_lower(kOperationColumnName);
319
320
1.50k
    _operation_idx = -1;
321
1.50k
    _row_id_idx = -1;
322
11.9k
    for (size_t i = 0; i < _vec_output_expr_ctxs.size(); ++i) {
323
10.4k
        std::string expr_name = doris::to_lower(_vec_output_expr_ctxs[i]->expr_name());
324
10.4k
        if (_operation_idx < 0 && expr_name == op_name) {
325
1.50k
            _operation_idx = static_cast<int>(i);
326
8.92k
        } else if (_row_id_idx < 0 && expr_name == row_id_name) {
327
1.50k
            _row_id_idx = static_cast<int>(i);
328
1.50k
        }
329
10.4k
    }
330
331
1.50k
    if (_operation_idx < 0) {
332
1
        return Status::InternalError("Iceberg merge sink missing operation column");
333
1
    }
334
1.50k
    if (_row_id_idx < 0) {
335
1
        return Status::InternalError("Iceberg merge sink missing row_id column");
336
1
    }
337
338
1.50k
    _data_column_indices.clear();
339
1.50k
    _table_output_expr_ctxs.clear();
340
11.9k
    for (size_t i = 0; i < _vec_output_expr_ctxs.size(); ++i) {
341
10.4k
        if (static_cast<int>(i) == _operation_idx || static_cast<int>(i) == _row_id_idx) {
342
3.01k
            continue;
343
3.01k
        }
344
7.42k
        _data_column_indices.push_back(static_cast<int>(i));
345
7.42k
        _table_output_expr_ctxs.emplace_back(_vec_output_expr_ctxs[i]);
346
7.42k
    }
347
348
1.50k
    return Status::OK();
349
1.50k
}
350
351
} // namespace doris