Coverage Report

Created: 2026-03-13 09:37

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/storage/partial_update_info.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 "storage/partial_update_info.h"
19
20
#include <gen_cpp/olap_file.pb.h>
21
22
#include <cstdint>
23
24
#include "common/consts.h"
25
#include "common/logging.h"
26
#include "core/assert_cast.h"
27
#include "core/block/block.h"
28
#include "core/data_type/data_type_number.h" // IWYU pragma: keep
29
#include "core/value/bitmap_value.h"
30
#include "storage/iterator/olap_data_convertor.h"
31
#include "storage/olap_common.h"
32
#include "storage/rowset/rowset.h"
33
#include "storage/rowset/rowset_writer_context.h"
34
#include "storage/segment/vertical_segment_writer.h"
35
#include "storage/tablet/base_tablet.h"
36
#include "storage/tablet/tablet_meta.h"
37
#include "storage/tablet/tablet_schema.h"
38
#include "storage/utils.h"
39
40
namespace doris {
41
#include "common/compile_check_begin.h"
42
Status PartialUpdateInfo::init(int64_t tablet_id, int64_t txn_id, const TabletSchema& tablet_schema,
43
                               UniqueKeyUpdateModePB unique_key_update_mode,
44
                               PartialUpdateNewRowPolicyPB policy,
45
                               const std::set<std::string>& partial_update_cols,
46
                               bool is_strict_mode_, int64_t timestamp_ms_, int32_t nano_seconds_,
47
                               const std::string& timezone_,
48
                               const std::string& auto_increment_column,
49
197k
                               int32_t sequence_map_col_uid, int64_t cur_max_version) {
50
197k
    partial_update_mode = unique_key_update_mode;
51
197k
    partial_update_new_key_policy = policy;
52
197k
    partial_update_input_columns = partial_update_cols;
53
197k
    max_version_in_flush_phase = cur_max_version;
54
197k
    sequence_map_col_unqiue_id = sequence_map_col_uid;
55
197k
    timestamp_ms = timestamp_ms_;
56
197k
    nano_seconds = nano_seconds_;
57
197k
    timezone = timezone_;
58
197k
    missing_cids.clear();
59
197k
    update_cids.clear();
60
61
197k
    if (partial_update_mode == UniqueKeyUpdateModePB::UPDATE_FIXED_COLUMNS) {
62
        // partial_update_cols should include all key columns
63
11.2k
        for (std::size_t i {0}; i < tablet_schema.num_key_columns(); i++) {
64
7.75k
            const auto key_col = tablet_schema.column(i);
65
7.75k
            if (!partial_update_cols.contains(key_col.name())) {
66
0
                auto msg = fmt::format(
67
0
                        "Unable to do partial update on shadow index's tablet, tablet_id={}, "
68
0
                        "txn_id={}. Missing key column {}.",
69
0
                        tablet_id, txn_id, key_col.name());
70
0
                LOG_WARNING(msg);
71
0
                return Status::Aborted<false>(msg);
72
0
            }
73
7.75k
        }
74
3.45k
    }
75
197k
    if (is_partial_update()) {
76
36.6k
        for (auto i = 0; i < tablet_schema.num_columns(); ++i) {
77
32.9k
            if (partial_update_mode == UniqueKeyUpdateModePB::UPDATE_FIXED_COLUMNS) {
78
29.0k
                auto tablet_column = tablet_schema.column(i);
79
29.0k
                if (!partial_update_input_columns.contains(tablet_column.name())) {
80
16.8k
                    missing_cids.emplace_back(i);
81
16.8k
                    if (!tablet_column.has_default_value() && !tablet_column.is_nullable() &&
82
16.8k
                        tablet_schema.auto_increment_column() != tablet_column.name()) {
83
1.24k
                        can_insert_new_rows_in_partial_update = false;
84
1.24k
                    }
85
16.8k
                } else {
86
12.1k
                    update_cids.emplace_back(i);
87
12.1k
                }
88
29.0k
                if (auto_increment_column == tablet_column.name()) {
89
240
                    is_schema_contains_auto_inc_column = true;
90
240
                }
91
29.0k
            } else {
92
                // in flexible partial update, missing cids is all non sort keys' cid
93
3.87k
                if (i >= tablet_schema.num_key_columns()) {
94
3.59k
                    missing_cids.emplace_back(i);
95
3.59k
                }
96
3.87k
            }
97
32.9k
        }
98
3.71k
        _generate_default_values_for_missing_cids(tablet_schema);
99
3.71k
    }
100
197k
    is_strict_mode = is_strict_mode_;
101
197k
    is_input_columns_contains_auto_inc_column =
102
197k
            is_fixed_partial_update() &&
103
197k
            partial_update_input_columns.contains(auto_increment_column);
104
197k
    return Status::OK();
105
197k
}
106
107
232
void PartialUpdateInfo::to_pb(PartialUpdateInfoPB* partial_update_info_pb) const {
108
232
    partial_update_info_pb->set_partial_update_mode(partial_update_mode);
109
232
    partial_update_info_pb->set_partial_update_new_key_policy(partial_update_new_key_policy);
110
232
    partial_update_info_pb->set_max_version_in_flush_phase(max_version_in_flush_phase);
111
1.53k
    for (const auto& col : partial_update_input_columns) {
112
1.53k
        partial_update_info_pb->add_partial_update_input_columns(col);
113
1.53k
    }
114
1.71k
    for (auto cid : missing_cids) {
115
1.71k
        partial_update_info_pb->add_missing_cids(cid);
116
1.71k
    }
117
1.53k
    for (auto cid : update_cids) {
118
1.53k
        partial_update_info_pb->add_update_cids(cid);
119
1.53k
    }
120
232
    partial_update_info_pb->set_can_insert_new_rows_in_partial_update(
121
232
            can_insert_new_rows_in_partial_update);
122
232
    partial_update_info_pb->set_is_strict_mode(is_strict_mode);
123
232
    partial_update_info_pb->set_timestamp_ms(timestamp_ms);
124
232
    partial_update_info_pb->set_nano_seconds(nano_seconds);
125
232
    partial_update_info_pb->set_timezone(timezone);
126
232
    partial_update_info_pb->set_is_input_columns_contains_auto_inc_column(
127
232
            is_input_columns_contains_auto_inc_column);
128
232
    partial_update_info_pb->set_is_schema_contains_auto_inc_column(
129
232
            is_schema_contains_auto_inc_column);
130
1.71k
    for (const auto& value : default_values) {
131
1.71k
        partial_update_info_pb->add_default_values(value);
132
1.71k
    }
133
232
}
134
135
0
void PartialUpdateInfo::from_pb(PartialUpdateInfoPB* partial_update_info_pb) {
136
0
    if (!partial_update_info_pb->has_partial_update_mode()) {
137
        // for backward compatibility
138
0
        if (partial_update_info_pb->is_partial_update()) {
139
0
            partial_update_mode = UniqueKeyUpdateModePB::UPDATE_FIXED_COLUMNS;
140
0
        } else {
141
0
            partial_update_mode = UniqueKeyUpdateModePB::UPSERT;
142
0
        }
143
0
    } else {
144
0
        partial_update_mode = partial_update_info_pb->partial_update_mode();
145
0
    }
146
0
    if (partial_update_info_pb->has_partial_update_new_key_policy()) {
147
0
        partial_update_new_key_policy = partial_update_info_pb->partial_update_new_key_policy();
148
0
    }
149
0
    max_version_in_flush_phase = partial_update_info_pb->has_max_version_in_flush_phase()
150
0
                                         ? partial_update_info_pb->max_version_in_flush_phase()
151
0
                                         : -1;
152
0
    partial_update_input_columns.clear();
153
0
    for (const auto& col : partial_update_info_pb->partial_update_input_columns()) {
154
0
        partial_update_input_columns.insert(col);
155
0
    }
156
0
    missing_cids.clear();
157
0
    for (auto cid : partial_update_info_pb->missing_cids()) {
158
0
        missing_cids.push_back(cid);
159
0
    }
160
0
    update_cids.clear();
161
0
    for (auto cid : partial_update_info_pb->update_cids()) {
162
0
        update_cids.push_back(cid);
163
0
    }
164
0
    can_insert_new_rows_in_partial_update =
165
0
            partial_update_info_pb->can_insert_new_rows_in_partial_update();
166
0
    is_strict_mode = partial_update_info_pb->is_strict_mode();
167
0
    timestamp_ms = partial_update_info_pb->timestamp_ms();
168
0
    timezone = partial_update_info_pb->timezone();
169
0
    is_input_columns_contains_auto_inc_column =
170
0
            partial_update_info_pb->is_input_columns_contains_auto_inc_column();
171
0
    is_schema_contains_auto_inc_column =
172
0
            partial_update_info_pb->is_schema_contains_auto_inc_column();
173
0
    if (partial_update_info_pb->has_nano_seconds()) {
174
0
        nano_seconds = partial_update_info_pb->nano_seconds();
175
0
    }
176
0
    default_values.clear();
177
0
    for (const auto& value : partial_update_info_pb->default_values()) {
178
0
        default_values.push_back(value);
179
0
    }
180
0
}
181
182
0
std::string PartialUpdateInfo::summary() const {
183
0
    std::string mode;
184
0
    switch (partial_update_mode) {
185
0
    case UniqueKeyUpdateModePB::UPSERT:
186
0
        mode = "upsert";
187
0
        break;
188
0
    case UniqueKeyUpdateModePB::UPDATE_FIXED_COLUMNS:
189
0
        mode = "fixed partial update";
190
0
        break;
191
0
    case UniqueKeyUpdateModePB::UPDATE_FLEXIBLE_COLUMNS:
192
0
        mode = "flexible partial update";
193
0
        break;
194
0
    }
195
0
    return fmt::format(
196
0
            "mode={}, update_cids={}, missing_cids={}, is_strict_mode={}, "
197
0
            "max_version_in_flush_phase={}",
198
0
            mode, update_cids.size(), missing_cids.size(), is_strict_mode,
199
0
            max_version_in_flush_phase);
200
0
}
201
202
Status PartialUpdateInfo::handle_new_key(const TabletSchema& tablet_schema,
203
                                         const std::function<std::string()>& line,
204
557
                                         BitmapValue* skip_bitmap) {
205
557
    switch (partial_update_new_key_policy) {
206
521
    case doris::PartialUpdateNewRowPolicyPB::APPEND: {
207
521
        if (partial_update_mode == UniqueKeyUpdateModePB::UPDATE_FIXED_COLUMNS) {
208
237
            if (!can_insert_new_rows_in_partial_update) {
209
10
                std::string error_column;
210
18
                for (auto cid : missing_cids) {
211
18
                    const TabletColumn& col = tablet_schema.column(cid);
212
18
                    if (!col.has_default_value() && !col.is_nullable() &&
213
18
                        !(tablet_schema.auto_increment_column() == col.name())) {
214
10
                        error_column = col.name();
215
10
                        break;
216
10
                    }
217
18
                }
218
10
                return Status::Error<ErrorCode::INVALID_SCHEMA, false>(
219
10
                        "the unmentioned column `{}` should have default value or be nullable "
220
10
                        "for newly inserted rows in non-strict mode partial update",
221
10
                        error_column);
222
10
            }
223
284
        } else if (partial_update_mode == UniqueKeyUpdateModePB::UPDATE_FLEXIBLE_COLUMNS) {
224
284
            DCHECK(skip_bitmap != nullptr);
225
284
            bool can_insert_new_row {true};
226
284
            std::string error_column;
227
4.03k
            for (auto cid : missing_cids) {
228
4.03k
                const TabletColumn& col = tablet_schema.column(cid);
229
4.03k
                if (skip_bitmap->contains(col.unique_id()) && !col.has_default_value() &&
230
4.03k
                    !col.is_nullable() && !col.is_auto_increment()) {
231
0
                    error_column = col.name();
232
0
                    can_insert_new_row = false;
233
0
                    break;
234
0
                }
235
4.03k
            }
236
284
            if (!can_insert_new_row) {
237
0
                return Status::Error<ErrorCode::INVALID_SCHEMA, false>(
238
0
                        "the unmentioned column `{}` should have default value or be "
239
0
                        "nullable for newly inserted rows in non-strict mode flexible partial "
240
0
                        "update",
241
0
                        error_column);
242
0
            }
243
284
        }
244
521
    } break;
245
511
    case doris::PartialUpdateNewRowPolicyPB::ERROR: {
246
36
        return Status::Error<ErrorCode::NEW_ROWS_IN_PARTIAL_UPDATE, false>(
247
36
                "Can't append new rows in partial update when partial_update_new_key_behavior is "
248
36
                "ERROR. Row with key=[{}] is not in table.",
249
36
                line());
250
521
    } break;
251
557
    }
252
512
    return Status::OK();
253
557
}
254
255
void PartialUpdateInfo::_generate_default_values_for_missing_cids(
256
3.71k
        const TabletSchema& tablet_schema) {
257
20.4k
    for (unsigned int cur_cid : missing_cids) {
258
20.4k
        const auto& column = tablet_schema.column(cur_cid);
259
20.4k
        if (column.has_default_value()) {
260
7.54k
            std::string default_value;
261
7.54k
            if (UNLIKELY((column.type() == FieldType::OLAP_FIELD_TYPE_DATETIMEV2 ||
262
7.54k
                          column.type() == FieldType::OLAP_FIELD_TYPE_TIMESTAMPTZ) &&
263
7.54k
                         to_lower(column.default_value()).find(to_lower("CURRENT_TIMESTAMP")) !=
264
7.54k
                                 std::string::npos)) {
265
133
                auto pos = to_lower(column.default_value()).find('(');
266
133
                if (pos == std::string::npos) {
267
64
                    DateV2Value<DateTimeV2ValueType> dtv;
268
64
                    dtv.from_unixtime(timestamp_ms / 1000, timezone);
269
64
                    default_value = dtv.to_string();
270
64
                    if (column.type() == FieldType::OLAP_FIELD_TYPE_TIMESTAMPTZ) {
271
0
                        default_value += timezone;
272
0
                    }
273
69
                } else {
274
69
                    int precision = std::stoi(column.default_value().substr(pos + 1));
275
69
                    DateV2Value<DateTimeV2ValueType> dtv;
276
69
                    dtv.from_unixtime(timestamp_ms / 1000, nano_seconds, timezone, precision);
277
69
                    default_value = dtv.to_string();
278
69
                    if (column.type() == FieldType::OLAP_FIELD_TYPE_TIMESTAMPTZ) {
279
0
                        default_value += timezone;
280
0
                    }
281
69
                }
282
7.41k
            } else if (UNLIKELY(column.type() == FieldType::OLAP_FIELD_TYPE_DATEV2 &&
283
7.41k
                                to_lower(column.default_value()).find(to_lower("CURRENT_DATE")) !=
284
7.41k
                                        std::string::npos)) {
285
81
                DateV2Value<DateV2ValueType> dv;
286
81
                dv.from_unixtime(timestamp_ms / 1000, timezone);
287
81
                default_value = dv.to_string();
288
7.33k
            } else if (UNLIKELY(column.type() == FieldType::OLAP_FIELD_TYPE_BITMAP &&
289
7.33k
                                to_lower(column.default_value()).find(to_lower("BITMAP_EMPTY")) !=
290
7.33k
                                        std::string::npos)) {
291
0
                BitmapValue v = BitmapValue {};
292
0
                default_value = v.to_string();
293
7.33k
            } else {
294
7.33k
                default_value = column.default_value();
295
7.33k
            }
296
7.54k
            default_values.emplace_back(default_value);
297
12.8k
        } else {
298
            // place an empty string here
299
12.8k
            default_values.emplace_back();
300
12.8k
        }
301
20.4k
    }
302
3.71k
    CHECK_EQ(missing_cids.size(), default_values.size());
303
3.71k
}
304
305
9
bool FixedReadPlan::empty() const {
306
9
    return plan.empty();
307
9
}
308
309
22.7k
void FixedReadPlan::prepare_to_read(const RowLocation& row_location, size_t pos) {
310
22.7k
    plan[row_location.rowset_id][row_location.segment_id].emplace_back(row_location.row_id, pos);
311
22.7k
}
312
313
// read columns by read plan
314
// read_index: ori_pos-> block_idx
315
Status FixedReadPlan::read_columns_by_plan(
316
        const TabletSchema& tablet_schema, std::vector<uint32_t> cids_to_read,
317
        const std::map<RowsetId, RowsetSharedPtr>& rsid_to_rowset, Block& block,
318
        std::map<uint32_t, uint32_t>* read_index, bool force_read_old_delete_signs,
319
1.34k
        const signed char* __restrict cur_delete_signs) const {
320
1.34k
    if (force_read_old_delete_signs) {
321
        // always read delete sign column from historical data
322
1.32k
        if (block.get_position_by_name(DELETE_SIGN) == -1) {
323
796
            auto del_col_cid = tablet_schema.field_index(DELETE_SIGN);
324
796
            cids_to_read.emplace_back(del_col_cid);
325
796
            block.swap(tablet_schema.create_block_by_cids(cids_to_read));
326
796
        }
327
1.32k
    }
328
1.34k
    bool has_row_column = tablet_schema.has_row_store_for_all_columns();
329
1.34k
    auto mutable_columns = block.mutate_columns();
330
1.34k
    uint32_t read_idx = 0;
331
1.34k
    for (const auto& [rowset_id, segment_row_mappings] : plan) {
332
622
        for (const auto& [segment_id, mappings] : segment_row_mappings) {
333
622
            auto rowset_iter = rsid_to_rowset.find(rowset_id);
334
622
            CHECK(rowset_iter != rsid_to_rowset.end());
335
622
            std::vector<uint32_t> rids;
336
22.7k
            for (auto [rid, pos] : mappings) {
337
22.7k
                if (cur_delete_signs && cur_delete_signs[pos]) {
338
3
                    continue;
339
3
                }
340
22.7k
                rids.emplace_back(rid);
341
22.7k
                (*read_index)[static_cast<uint32_t>(pos)] = read_idx++;
342
22.7k
            }
343
622
            if (has_row_column) {
344
177
                auto st = BaseTablet::fetch_value_through_row_column(
345
177
                        rowset_iter->second, tablet_schema, segment_id, rids, cids_to_read, block);
346
177
                if (!st.ok()) {
347
0
                    LOG(WARNING) << "failed to fetch value through row column";
348
0
                    return st;
349
0
                }
350
177
                continue;
351
177
            }
352
2.57k
            for (size_t cid = 0; cid < mutable_columns.size(); ++cid) {
353
2.13k
                TabletColumn tablet_column = tablet_schema.column(cids_to_read[cid]);
354
2.13k
                auto st = doris::BaseTablet::fetch_value_by_rowids(
355
2.13k
                        rowset_iter->second, segment_id, rids, tablet_column, mutable_columns[cid]);
356
                // set read value to output block
357
2.13k
                if (!st.ok()) {
358
0
                    LOG(WARNING) << "failed to fetch value";
359
0
                    return st;
360
0
                }
361
2.13k
            }
362
445
        }
363
621
    }
364
1.34k
    block.set_columns(std::move(mutable_columns));
365
1.34k
    return Status::OK();
366
1.34k
}
367
368
Status FixedReadPlan::fill_missing_columns(
369
        RowsetWriterContext* rowset_ctx, const std::map<RowsetId, RowsetSharedPtr>& rsid_to_rowset,
370
        const TabletSchema& tablet_schema, Block& full_block,
371
        const std::vector<bool>& use_default_or_null_flag, bool has_default_or_nullable,
372
1.30k
        uint32_t segment_start_pos, const Block* block) const {
373
1.30k
    auto mutable_full_columns = full_block.mutate_columns();
374
    // create old value columns
375
1.30k
    const auto& missing_cids = rowset_ctx->partial_update_info->missing_cids;
376
1.30k
    bool have_input_seq_column = false;
377
1.30k
    if (tablet_schema.has_sequence_col()) {
378
224
        const std::vector<uint32_t>& including_cids = rowset_ctx->partial_update_info->update_cids;
379
224
        have_input_seq_column =
380
224
                (std::find(including_cids.cbegin(), including_cids.cend(),
381
224
                           tablet_schema.sequence_col_idx()) != including_cids.cend());
382
224
    }
383
384
1.30k
    auto old_value_block = tablet_schema.create_block_by_cids(missing_cids);
385
1.30k
    CHECK_EQ(missing_cids.size(), old_value_block.columns());
386
387
    // segment pos to write -> rowid to read in old_value_block
388
1.30k
    std::map<uint32_t, uint32_t> read_index;
389
1.30k
    RETURN_IF_ERROR(read_columns_by_plan(tablet_schema, missing_cids, rsid_to_rowset,
390
1.30k
                                         old_value_block, &read_index, true, nullptr));
391
392
1.30k
    const auto* old_delete_signs = BaseTablet::get_delete_sign_column_data(old_value_block);
393
1.30k
    if (old_delete_signs == nullptr) {
394
0
        return Status::InternalError("old delete signs column not found, block: {}",
395
0
                                     old_value_block.dump_structure());
396
0
    }
397
    // build default value columns
398
1.30k
    auto default_value_block = old_value_block.clone_empty();
399
1.30k
    RETURN_IF_ERROR(BaseTablet::generate_default_value_block(
400
1.30k
            tablet_schema, missing_cids, rowset_ctx->partial_update_info->default_values,
401
1.30k
            old_value_block, default_value_block));
402
1.30k
    auto mutable_default_value_columns = default_value_block.mutate_columns();
403
404
    // fill all missing value from mutable_old_columns, need to consider default value and null value
405
25.6k
    for (auto idx = 0; idx < use_default_or_null_flag.size(); idx++) {
406
24.2k
        auto segment_pos = idx + segment_start_pos;
407
24.2k
        auto pos_in_old_block = read_index[segment_pos];
408
409
115k
        for (auto i = 0; i < missing_cids.size(); ++i) {
410
            // if the column has default value, fill it with default value
411
            // otherwise, if the column is nullable, fill it with null value
412
91.6k
            const auto& tablet_column = tablet_schema.column(missing_cids[i]);
413
91.6k
            auto& missing_col = mutable_full_columns[missing_cids[i]];
414
415
91.6k
            bool should_use_default = use_default_or_null_flag[idx];
416
91.6k
            if (!should_use_default) {
417
79.0k
                bool old_row_delete_sign =
418
79.0k
                        (old_delete_signs != nullptr && old_delete_signs[pos_in_old_block] != 0);
419
79.0k
                if (old_row_delete_sign) {
420
151
                    if (!tablet_schema.has_sequence_col()) {
421
70
                        should_use_default = true;
422
81
                    } else if (have_input_seq_column || (!tablet_column.is_seqeunce_col())) {
423
                        // to keep the sequence column value not decreasing, we should read values of seq column
424
                        // from old rows even if the old row is deleted when the input don't specify the sequence column, otherwise
425
                        // it may cause the merge-on-read based compaction to produce incorrect results
426
77
                        should_use_default = true;
427
77
                    }
428
151
                }
429
79.0k
            }
430
431
91.6k
            if (should_use_default) {
432
12.8k
                if (tablet_column.has_default_value()) {
433
2.76k
                    missing_col->insert_from(*mutable_default_value_columns[i], 0);
434
10.0k
                } else if (tablet_column.is_nullable()) {
435
8.52k
                    auto* nullable_column = assert_cast<ColumnNullable*>(missing_col.get());
436
8.52k
                    nullable_column->insert_many_defaults(1);
437
8.52k
                } else if (tablet_schema.auto_increment_column() == tablet_column.name()) {
438
41
                    const auto& column =
439
41
                            *DORIS_TRY(rowset_ctx->tablet_schema->column(tablet_column.name()));
440
41
                    DCHECK(column.type() == FieldType::OLAP_FIELD_TYPE_BIGINT);
441
41
                    auto* auto_inc_column = assert_cast<ColumnInt64*>(missing_col.get());
442
41
                    int pos = block->get_position_by_name(BeConsts::PARTIAL_UPDATE_AUTO_INC_COL);
443
41
                    if (pos == -1) {
444
0
                        return Status::InternalError("auto increment column not found in block {}",
445
0
                                                     block->dump_structure());
446
0
                    }
447
41
                    auto_inc_column->insert_from(*block->get_by_position(pos).column.get(), idx);
448
1.49k
                } else {
449
                    // If the control flow reaches this branch, the column neither has default value
450
                    // nor is nullable. It means that the row's delete sign is marked, and the value
451
                    // columns are useless and won't be read. So we can just put arbitary values in the cells
452
1.49k
                    missing_col->insert(tablet_column.get_vec_type()->get_default());
453
1.49k
                }
454
78.8k
            } else {
455
78.8k
                missing_col->insert_from(*old_value_block.get_by_position(i).column,
456
78.8k
                                         pos_in_old_block);
457
78.8k
            }
458
91.6k
        }
459
24.2k
    }
460
1.30k
    full_block.set_columns(std::move(mutable_full_columns));
461
1.30k
    return Status::OK();
462
1.30k
}
463
464
void FlexibleReadPlan::prepare_to_read(const RowLocation& row_location, size_t pos,
465
1.79k
                                       const BitmapValue& skip_bitmap) {
466
1.79k
    if (!use_row_store) {
467
3.90k
        for (uint64_t col_uid : skip_bitmap) {
468
3.90k
            plan[row_location.rowset_id][row_location.segment_id][static_cast<uint32_t>(col_uid)]
469
3.90k
                    .emplace_back(row_location.row_id, pos);
470
3.90k
        }
471
932
    } else {
472
859
        row_store_plan[row_location.rowset_id][row_location.segment_id].emplace_back(
473
859
                row_location.row_id, pos);
474
859
    }
475
1.79k
}
476
477
Status FlexibleReadPlan::read_columns_by_plan(
478
        const TabletSchema& tablet_schema,
479
        const std::map<RowsetId, RowsetSharedPtr>& rsid_to_rowset, Block& old_value_block,
480
145
        std::map<uint32_t, std::map<uint32_t, uint32_t>>* read_index) const {
481
145
    auto mutable_columns = old_value_block.mutate_columns();
482
483
    // cid -> next rid to fill in block
484
145
    std::map<uint32_t, uint32_t> next_read_idx;
485
2.16k
    for (uint32_t cid {0}; cid < tablet_schema.num_columns(); cid++) {
486
2.02k
        next_read_idx[cid] = 0;
487
2.02k
    }
488
489
180
    for (const auto& [rowset_id, segment_mappings] : plan) {
490
180
        for (const auto& [segment_id, uid_mappings] : segment_mappings) {
491
1.48k
            for (const auto& [col_uid, mappings] : uid_mappings) {
492
1.48k
                auto rowset_iter = rsid_to_rowset.find(rowset_id);
493
1.48k
                CHECK(rowset_iter != rsid_to_rowset.end());
494
1.48k
                auto cid = tablet_schema.field_index(col_uid);
495
1.48k
                DCHECK_NE(cid, -1);
496
1.48k
                DCHECK_GE(cid, tablet_schema.num_key_columns());
497
1.48k
                std::vector<uint32_t> rids;
498
3.91k
                for (auto [rid, pos] : mappings) {
499
3.91k
                    rids.emplace_back(rid);
500
3.91k
                    (*read_index)[cid][static_cast<uint32_t>(pos)] = next_read_idx[cid]++;
501
3.91k
                }
502
503
1.48k
                TabletColumn tablet_column = tablet_schema.column(cid);
504
1.48k
                auto idx = cid - tablet_schema.num_key_columns();
505
1.48k
                RETURN_IF_ERROR(doris::BaseTablet::fetch_value_by_rowids(
506
1.48k
                        rowset_iter->second, segment_id, rids, tablet_column,
507
1.48k
                        mutable_columns[idx]));
508
1.48k
            }
509
180
        }
510
180
    }
511
    // !!!ATTENTION!!!: columns in block may have different size because every row has different columns to update
512
145
    old_value_block.set_columns(std::move(mutable_columns));
513
145
    return Status::OK();
514
145
}
515
516
Status FlexibleReadPlan::read_columns_by_plan(
517
        const TabletSchema& tablet_schema, const std::vector<uint32_t>& cids_to_read,
518
        const std::map<RowsetId, RowsetSharedPtr>& rsid_to_rowset, Block& old_value_block,
519
119
        std::map<uint32_t, uint32_t>* read_index) const {
520
119
    DCHECK(use_row_store);
521
119
    uint32_t read_idx = 0;
522
158
    for (const auto& [rowset_id, segment_row_mappings] : row_store_plan) {
523
158
        for (const auto& [segment_id, mappings] : segment_row_mappings) {
524
158
            auto rowset_iter = rsid_to_rowset.find(rowset_id);
525
158
            CHECK(rowset_iter != rsid_to_rowset.end());
526
158
            std::vector<uint32_t> rids;
527
859
            for (auto [rid, pos] : mappings) {
528
859
                rids.emplace_back(rid);
529
859
                (*read_index)[static_cast<uint32_t>(pos)] = read_idx++;
530
859
            }
531
158
            auto st = BaseTablet::fetch_value_through_row_column(rowset_iter->second, tablet_schema,
532
158
                                                                 segment_id, rids, cids_to_read,
533
158
                                                                 old_value_block);
534
158
            if (!st.ok()) {
535
0
                LOG(WARNING) << "failed to fetch value through row column";
536
0
                return st;
537
0
            }
538
158
        }
539
158
    }
540
119
    return Status::OK();
541
119
}
542
543
Status FlexibleReadPlan::fill_non_primary_key_columns(
544
        RowsetWriterContext* rowset_ctx, const std::map<RowsetId, RowsetSharedPtr>& rsid_to_rowset,
545
        const TabletSchema& tablet_schema, Block& full_block,
546
        const std::vector<bool>& use_default_or_null_flag, bool has_default_or_nullable,
547
        uint32_t segment_start_pos, uint32_t block_start_pos, const Block* block,
548
264
        std::vector<BitmapValue>* skip_bitmaps) const {
549
264
    auto mutable_full_columns = full_block.mutate_columns();
550
551
    // missing_cids are all non sort key columns' cids
552
264
    const auto& non_sort_key_cids = rowset_ctx->partial_update_info->missing_cids;
553
264
    auto old_value_block = tablet_schema.create_block_by_cids(non_sort_key_cids);
554
264
    CHECK_EQ(non_sort_key_cids.size(), old_value_block.columns());
555
556
264
    if (!use_row_store) {
557
145
        RETURN_IF_ERROR(fill_non_primary_key_columns_for_column_store(
558
145
                rowset_ctx, rsid_to_rowset, tablet_schema, non_sort_key_cids, old_value_block,
559
145
                mutable_full_columns, use_default_or_null_flag, has_default_or_nullable,
560
145
                segment_start_pos, block_start_pos, block, skip_bitmaps));
561
145
    } else {
562
119
        RETURN_IF_ERROR(fill_non_primary_key_columns_for_row_store(
563
119
                rowset_ctx, rsid_to_rowset, tablet_schema, non_sort_key_cids, old_value_block,
564
119
                mutable_full_columns, use_default_or_null_flag, has_default_or_nullable,
565
119
                segment_start_pos, block_start_pos, block, skip_bitmaps));
566
119
    }
567
264
    full_block.set_columns(std::move(mutable_full_columns));
568
264
    return Status::OK();
569
264
}
570
571
static void fill_non_primary_key_cell_for_column_store(
572
        const TabletColumn& tablet_column, uint32_t cid, MutableColumnPtr& new_col,
573
        const IColumn& default_value_col, const IColumn& old_value_col, const IColumn& cur_col,
574
        std::size_t block_pos, uint32_t segment_pos, bool skipped, bool row_has_sequence_col,
575
        bool use_default, const signed char* delete_sign_column_data,
576
        const TabletSchema& tablet_schema,
577
        std::map<uint32_t, std::map<uint32_t, uint32_t>>& read_index,
578
16.9k
        const PartialUpdateInfo* info) {
579
16.9k
    if (skipped) {
580
4.78k
        DCHECK(cid != tablet_schema.skip_bitmap_col_idx());
581
4.78k
        DCHECK(cid != tablet_schema.version_col_idx());
582
4.78k
        DCHECK(!tablet_column.is_row_store_column());
583
584
4.78k
        if (!use_default) {
585
3.91k
            if (delete_sign_column_data != nullptr) {
586
3.91k
                bool old_row_delete_sign = false;
587
3.91k
                if (auto it = read_index[tablet_schema.delete_sign_idx()].find(segment_pos);
588
3.91k
                    it != read_index[tablet_schema.delete_sign_idx()].end()) {
589
3.79k
                    old_row_delete_sign = (delete_sign_column_data[it->second] != 0);
590
3.79k
                }
591
592
3.91k
                if (old_row_delete_sign) {
593
7
                    if (!tablet_schema.has_sequence_col()) {
594
2
                        use_default = true;
595
5
                    } else if (row_has_sequence_col ||
596
5
                               (!tablet_column.is_seqeunce_col() &&
597
5
                                (tablet_column.unique_id() != info->sequence_map_col_uid()))) {
598
                        // to keep the sequence column value not decreasing, we should read values of seq column(and seq map column)
599
                        // from old rows even if the old row is deleted when the input don't specify the sequence column, otherwise
600
                        // it may cause the merge-on-read based compaction to produce incorrect results
601
3
                        use_default = true;
602
3
                    }
603
7
                }
604
3.91k
            }
605
3.91k
        }
606
4.78k
        if (!use_default && tablet_column.is_on_update_current_timestamp()) {
607
6
            use_default = true;
608
6
        }
609
4.78k
        if (use_default) {
610
873
            if (tablet_column.has_default_value()) {
611
424
                new_col->insert_from(default_value_col, 0);
612
449
            } else if (tablet_column.is_nullable()) {
613
426
                assert_cast<ColumnNullable*, TypeCheckOnRelease::DISABLE>(new_col.get())
614
426
                        ->insert_many_defaults(1);
615
426
            } else if (tablet_column.is_auto_increment()) {
616
                // In flexible partial update, the skip bitmap indicates whether a cell
617
                // is specified in the original load, so the generated auto-increment value is filled
618
                // in current block in place if needed rather than using a seperate column to
619
                // store the generated auto-increment value in fixed partial update
620
2
                new_col->insert_from(cur_col, block_pos);
621
21
            } else {
622
21
                new_col->insert(tablet_column.get_vec_type()->get_default());
623
21
            }
624
3.90k
        } else {
625
3.90k
            auto pos_in_old_block = read_index.at(cid).at(segment_pos);
626
3.90k
            new_col->insert_from(old_value_col, pos_in_old_block);
627
3.90k
        }
628
12.1k
    } else {
629
12.1k
        new_col->insert_from(cur_col, block_pos);
630
12.1k
    }
631
16.9k
}
632
633
Status FlexibleReadPlan::fill_non_primary_key_columns_for_column_store(
634
        RowsetWriterContext* rowset_ctx, const std::map<RowsetId, RowsetSharedPtr>& rsid_to_rowset,
635
        const TabletSchema& tablet_schema, const std::vector<uint32_t>& non_sort_key_cids,
636
        Block& old_value_block, MutableColumns& mutable_full_columns,
637
        const std::vector<bool>& use_default_or_null_flag, bool has_default_or_nullable,
638
        uint32_t segment_start_pos, uint32_t block_start_pos, const Block* block,
639
145
        std::vector<BitmapValue>* skip_bitmaps) const {
640
145
    auto* info = rowset_ctx->partial_update_info.get();
641
145
    int32_t seq_col_unique_id = -1;
642
145
    if (tablet_schema.has_sequence_col()) {
643
22
        seq_col_unique_id = tablet_schema.column(tablet_schema.sequence_col_idx()).unique_id();
644
22
    }
645
    // cid -> segment pos to write -> rowid to read in old_value_block
646
145
    std::map<uint32_t, std::map<uint32_t, uint32_t>> read_index;
647
145
    RETURN_IF_ERROR(
648
145
            read_columns_by_plan(tablet_schema, rsid_to_rowset, old_value_block, &read_index));
649
    // !!!ATTENTION!!!: columns in old_value_block may have different size because every row has different columns to update
650
651
145
    const auto* delete_sign_column_data = BaseTablet::get_delete_sign_column_data(old_value_block);
652
    // build default value columns
653
145
    auto default_value_block = old_value_block.clone_empty();
654
145
    if (has_default_or_nullable || delete_sign_column_data != nullptr) {
655
145
        RETURN_IF_ERROR(BaseTablet::generate_default_value_block(
656
145
                tablet_schema, non_sort_key_cids, info->default_values, old_value_block,
657
145
                default_value_block));
658
145
    }
659
660
    // fill all non sort key columns from mutable_old_columns, need to consider default value and null value
661
2.00k
    for (std::size_t i {0}; i < non_sort_key_cids.size(); i++) {
662
1.86k
        auto cid = non_sort_key_cids[i];
663
1.86k
        const auto& tablet_column = tablet_schema.column(cid);
664
1.86k
        auto col_uid = tablet_column.unique_id();
665
18.8k
        for (auto idx = 0; idx < use_default_or_null_flag.size(); idx++) {
666
16.9k
            auto segment_pos = segment_start_pos + idx;
667
16.9k
            auto block_pos = block_start_pos + idx;
668
669
16.9k
            fill_non_primary_key_cell_for_column_store(
670
16.9k
                    tablet_column, cid, mutable_full_columns[cid],
671
16.9k
                    *default_value_block.get_by_position(i).column,
672
16.9k
                    *old_value_block.get_by_position(i).column, *block->get_by_position(cid).column,
673
16.9k
                    block_pos, segment_pos, skip_bitmaps->at(block_pos).contains(col_uid),
674
16.9k
                    tablet_schema.has_sequence_col()
675
16.9k
                            ? !skip_bitmaps->at(block_pos).contains(seq_col_unique_id)
676
16.9k
                            : false,
677
16.9k
                    use_default_or_null_flag[idx], delete_sign_column_data, tablet_schema,
678
16.9k
                    read_index, info);
679
16.9k
        }
680
1.86k
    }
681
145
    return Status::OK();
682
145
}
683
684
static void fill_non_primary_key_cell_for_row_store(
685
        const TabletColumn& tablet_column, uint32_t cid, MutableColumnPtr& new_col,
686
        const IColumn& default_value_col, const IColumn& old_value_col, const IColumn& cur_col,
687
        std::size_t block_pos, bool skipped, bool row_has_sequence_col, bool use_default,
688
        const signed char* delete_sign_column_data, uint32_t pos_in_old_block,
689
16.5k
        const TabletSchema& tablet_schema, const PartialUpdateInfo* info) {
690
16.5k
    if (skipped) {
691
4.16k
        DCHECK(cid != tablet_schema.skip_bitmap_col_idx());
692
4.16k
        DCHECK(cid != tablet_schema.version_col_idx());
693
4.16k
        DCHECK(!tablet_column.is_row_store_column());
694
4.16k
        if (!use_default) {
695
3.64k
            if (delete_sign_column_data != nullptr) {
696
3.64k
                bool old_row_delete_sign = (delete_sign_column_data[pos_in_old_block] != 0);
697
3.64k
                if (old_row_delete_sign) {
698
7
                    if (!tablet_schema.has_sequence_col()) {
699
2
                        use_default = true;
700
5
                    } else if (row_has_sequence_col ||
701
5
                               (!tablet_column.is_seqeunce_col() &&
702
5
                                (tablet_column.unique_id() != info->sequence_map_col_uid()))) {
703
                        // to keep the sequence column value not decreasing, we should read values of seq column(and seq map column)
704
                        // from old rows even if the old row is deleted when the input don't specify the sequence column, otherwise
705
                        // it may cause the merge-on-read based compaction to produce incorrect results
706
3
                        use_default = true;
707
3
                    }
708
7
                }
709
3.64k
            }
710
3.64k
        }
711
712
4.16k
        if (!use_default && tablet_column.is_on_update_current_timestamp()) {
713
6
            use_default = true;
714
6
        }
715
4.16k
        if (use_default) {
716
529
            if (tablet_column.has_default_value()) {
717
175
                new_col->insert_from(default_value_col, 0);
718
354
            } else if (tablet_column.is_nullable()) {
719
350
                assert_cast<ColumnNullable*, TypeCheckOnRelease::DISABLE>(new_col.get())
720
350
                        ->insert_many_defaults(1);
721
350
            } else if (tablet_column.is_auto_increment()) {
722
                // In flexible partial update, the skip bitmap indicates whether a cell
723
                // is specified in the original load, so the generated auto-increment value is filled
724
                // in current block in place if needed rather than using a seperate column to
725
                // store the generated auto-increment value in fixed partial update
726
2
                new_col->insert_from(cur_col, block_pos);
727
2
            } else {
728
2
                new_col->insert(tablet_column.get_vec_type()->get_default());
729
2
            }
730
3.63k
        } else {
731
3.63k
            new_col->insert_from(old_value_col, pos_in_old_block);
732
3.63k
        }
733
12.3k
    } else {
734
12.3k
        new_col->insert_from(cur_col, block_pos);
735
12.3k
    }
736
16.5k
}
737
738
Status FlexibleReadPlan::fill_non_primary_key_columns_for_row_store(
739
        RowsetWriterContext* rowset_ctx, const std::map<RowsetId, RowsetSharedPtr>& rsid_to_rowset,
740
        const TabletSchema& tablet_schema, const std::vector<uint32_t>& non_sort_key_cids,
741
        Block& old_value_block, MutableColumns& mutable_full_columns,
742
        const std::vector<bool>& use_default_or_null_flag, bool has_default_or_nullable,
743
        uint32_t segment_start_pos, uint32_t block_start_pos, const Block* block,
744
119
        std::vector<BitmapValue>* skip_bitmaps) const {
745
119
    auto* info = rowset_ctx->partial_update_info.get();
746
119
    int32_t seq_col_unique_id = -1;
747
119
    if (tablet_schema.has_sequence_col()) {
748
3
        seq_col_unique_id = tablet_schema.column(tablet_schema.sequence_col_idx()).unique_id();
749
3
    }
750
    // segment pos to write -> rowid to read in old_value_block
751
119
    std::map<uint32_t, uint32_t> read_index;
752
119
    RETURN_IF_ERROR(read_columns_by_plan(tablet_schema, non_sort_key_cids, rsid_to_rowset,
753
119
                                         old_value_block, &read_index));
754
755
119
    const auto* delete_sign_column_data = BaseTablet::get_delete_sign_column_data(old_value_block);
756
    // build default value columns
757
119
    auto default_value_block = old_value_block.clone_empty();
758
119
    if (has_default_or_nullable || delete_sign_column_data != nullptr) {
759
119
        RETURN_IF_ERROR(BaseTablet::generate_default_value_block(
760
119
                tablet_schema, non_sort_key_cids, info->default_values, old_value_block,
761
119
                default_value_block));
762
119
    }
763
764
    // fill all non sort key columns from mutable_old_columns, need to consider default value and null value
765
1.87k
    for (std::size_t i {0}; i < non_sort_key_cids.size(); i++) {
766
1.75k
        auto cid = non_sort_key_cids[i];
767
1.75k
        const auto& tablet_column = tablet_schema.column(cid);
768
1.75k
        auto col_uid = tablet_column.unique_id();
769
18.2k
        for (auto idx = 0; idx < use_default_or_null_flag.size(); idx++) {
770
16.5k
            auto segment_pos = segment_start_pos + idx;
771
16.5k
            auto block_pos = block_start_pos + idx;
772
16.5k
            auto pos_in_old_block = read_index[segment_pos];
773
774
16.5k
            fill_non_primary_key_cell_for_row_store(
775
16.5k
                    tablet_column, cid, mutable_full_columns[cid],
776
16.5k
                    *default_value_block.get_by_position(i).column,
777
16.5k
                    *old_value_block.get_by_position(i).column, *block->get_by_position(cid).column,
778
16.5k
                    block_pos, skip_bitmaps->at(block_pos).contains(col_uid),
779
16.5k
                    tablet_schema.has_sequence_col()
780
16.5k
                            ? !skip_bitmaps->at(block_pos).contains(seq_col_unique_id)
781
16.5k
                            : false,
782
16.5k
                    use_default_or_null_flag[idx], delete_sign_column_data, pos_in_old_block,
783
16.5k
                    tablet_schema, info);
784
16.5k
        }
785
1.75k
    }
786
119
    return Status::OK();
787
119
}
788
789
BlockAggregator::BlockAggregator(segment_v2::VerticalSegmentWriter& vertical_segment_writer)
790
58.5k
        : _writer(vertical_segment_writer), _tablet_schema(*_writer._tablet_schema) {}
791
792
void BlockAggregator::merge_one_row(MutableBlock& dst_block, Block* src_block, int rid,
793
153
                                    BitmapValue& skip_bitmap) {
794
1.57k
    for (size_t cid {_tablet_schema.num_key_columns()}; cid < _tablet_schema.num_columns(); cid++) {
795
1.42k
        if (cid == _tablet_schema.skip_bitmap_col_idx()) {
796
153
            auto& cur_skip_bitmap =
797
153
                    assert_cast<ColumnBitmap*>(dst_block.mutable_columns()[cid].get())
798
153
                            ->get_data()
799
153
                            .back();
800
153
            const auto& new_row_skip_bitmap =
801
153
                    assert_cast<ColumnBitmap*>(
802
153
                            src_block->get_by_position(cid).column->assume_mutable().get())
803
153
                            ->get_data()[rid];
804
153
            cur_skip_bitmap &= new_row_skip_bitmap;
805
153
            continue;
806
153
        }
807
1.27k
        if (!skip_bitmap.contains(_tablet_schema.column(cid).unique_id())) {
808
460
            dst_block.mutable_columns()[cid]->pop_back(1);
809
460
            dst_block.mutable_columns()[cid]->insert_from(*src_block->get_by_position(cid).column,
810
460
                                                          rid);
811
460
        }
812
1.27k
    }
813
153
    VLOG_DEBUG << fmt::format("merge a row, after merge, output_block.rows()={}, state: {}",
814
0
                              dst_block.rows(), _state.to_string());
815
153
}
816
817
252
void BlockAggregator::append_one_row(MutableBlock& dst_block, Block* src_block, int rid) {
818
252
    dst_block.add_row(src_block, rid);
819
252
    _state.rows++;
820
252
    VLOG_DEBUG << fmt::format("append a new row, after append, output_block.rows()={}, state: {}",
821
0
                              dst_block.rows(), _state.to_string());
822
252
}
823
824
123
void BlockAggregator::remove_last_n_rows(MutableBlock& dst_block, int n) {
825
123
    if (n > 0) {
826
968
        for (size_t cid {0}; cid < _tablet_schema.num_columns(); cid++) {
827
880
            DCHECK_GE(dst_block.mutable_columns()[cid]->size(), n);
828
880
            dst_block.mutable_columns()[cid]->pop_back(n);
829
880
        }
830
88
    }
831
123
}
832
833
void BlockAggregator::append_or_merge_row(MutableBlock& dst_block, Block* src_block, int rid,
834
405
                                          BitmapValue& skip_bitmap, bool have_delete_sign) {
835
405
    if (have_delete_sign) {
836
        // remove all the previous batched rows
837
123
        remove_last_n_rows(dst_block, _state.rows);
838
123
        _state.rows = 0;
839
123
        _state.has_row_with_delete_sign = true;
840
841
123
        append_one_row(dst_block, src_block, rid);
842
282
    } else {
843
282
        if (_state.should_merge()) {
844
153
            merge_one_row(dst_block, src_block, rid, skip_bitmap);
845
153
        } else {
846
129
            append_one_row(dst_block, src_block, rid);
847
129
        }
848
282
    }
849
405
};
850
851
Status BlockAggregator::aggregate_rows(
852
        MutableBlock& output_block, Block* block, int start, int end, std::string key,
853
        std::vector<BitmapValue>* skip_bitmaps, const signed char* delete_signs,
854
        IOlapColumnDataAccessor* seq_column, const std::vector<RowsetSharedPtr>& specified_rowsets,
855
133
        std::vector<std::unique_ptr<SegmentCacheHandle>>& segment_caches) {
856
133
    VLOG_DEBUG << fmt::format("merge rows in range=[{}-{})", start, end);
857
133
    if (end - start == 1) {
858
34
        output_block.add_row(block, start);
859
34
        VLOG_DEBUG << fmt::format("append a row directly, rid={}", start);
860
34
        return Status::OK();
861
34
    }
862
863
99
    auto seq_col_unique_id = _tablet_schema.column(_tablet_schema.sequence_col_idx()).unique_id();
864
99
    auto delete_sign_col_unique_id =
865
99
            _tablet_schema.column(_tablet_schema.delete_sign_idx()).unique_id();
866
867
99
    _state.reset();
868
869
99
    RowLocation loc;
870
99
    RowsetSharedPtr rowset;
871
99
    std::string previous_encoded_seq_value {};
872
99
    Status st = _writer._tablet->lookup_row_key(
873
99
            key, &_tablet_schema, false, specified_rowsets, &loc, _writer._mow_context->max_version,
874
99
            segment_caches, &rowset, true, &previous_encoded_seq_value);
875
99
    int pos = start;
876
99
    bool is_expected_st = (st.is<ErrorCode::KEY_NOT_FOUND>() || st.ok());
877
99
    DCHECK(is_expected_st || st.is<ErrorCode::MEM_LIMIT_EXCEEDED>())
878
0
            << "[BlockAggregator::aggregate_rows] unexpected error status while lookup_row_key:"
879
0
            << st;
880
99
    if (!is_expected_st) {
881
0
        return st;
882
0
    }
883
884
99
    std::string cur_seq_val;
885
99
    if (st.ok()) {
886
97
        for (pos = start; pos < end; pos++) {
887
97
            auto& skip_bitmap = skip_bitmaps->at(pos);
888
97
            bool row_has_sequence_col = (!skip_bitmap.contains(seq_col_unique_id));
889
            // Discard all the rows whose seq value is smaller than previous_encoded_seq_value.
890
97
            if (row_has_sequence_col) {
891
56
                std::string seq_val {};
892
56
                _writer._encode_seq_column(seq_column, pos, &seq_val);
893
56
                if (Slice {seq_val}.compare(Slice {previous_encoded_seq_value}) < 0) {
894
19
                    continue;
895
19
                }
896
37
                cur_seq_val = std::move(seq_val);
897
37
                break;
898
56
            }
899
41
            cur_seq_val = std::move(previous_encoded_seq_value);
900
41
            break;
901
97
        }
902
78
    } else {
903
21
        pos = start;
904
21
        auto& skip_bitmap = skip_bitmaps->at(pos);
905
21
        bool row_has_sequence_col = (!skip_bitmap.contains(seq_col_unique_id));
906
21
        if (row_has_sequence_col) {
907
15
            std::string seq_val {};
908
            // for rows that don't specify seqeunce col, seq_val will be encoded to minial value
909
15
            _writer._encode_seq_column(seq_column, pos, &seq_val);
910
15
            cur_seq_val = std::move(seq_val);
911
15
        } else {
912
6
            cur_seq_val.clear();
913
6
            RETURN_IF_ERROR(_writer._generate_encoded_default_seq_value(
914
6
                    _tablet_schema, *_writer._opts.rowset_ctx->partial_update_info, &cur_seq_val));
915
6
        }
916
21
    }
917
918
562
    for (int rid {pos}; rid < end; rid++) {
919
463
        auto& skip_bitmap = skip_bitmaps->at(rid);
920
463
        bool row_has_sequence_col = (!skip_bitmap.contains(seq_col_unique_id));
921
463
        bool have_delete_sign =
922
463
                (!skip_bitmap.contains(delete_sign_col_unique_id) && delete_signs[rid] != 0);
923
463
        if (!row_has_sequence_col) {
924
193
            append_or_merge_row(output_block, block, rid, skip_bitmap, have_delete_sign);
925
270
        } else {
926
270
            std::string seq_val {};
927
270
            _writer._encode_seq_column(seq_column, rid, &seq_val);
928
270
            if (Slice {seq_val}.compare(Slice {cur_seq_val}) >= 0) {
929
212
                append_or_merge_row(output_block, block, rid, skip_bitmap, have_delete_sign);
930
212
                cur_seq_val = std::move(seq_val);
931
212
            } else {
932
58
                VLOG_DEBUG << fmt::format(
933
0
                        "skip rid={} becasue its seq value is lower than the previous", rid);
934
58
            }
935
270
        }
936
463
    }
937
99
    return Status::OK();
938
99
};
939
940
Status BlockAggregator::aggregate_for_sequence_column(
941
        Block* block, int num_rows, const std::vector<IOlapColumnDataAccessor*>& key_columns,
942
        IOlapColumnDataAccessor* seq_column, const std::vector<RowsetSharedPtr>& specified_rowsets,
943
25
        std::vector<std::unique_ptr<SegmentCacheHandle>>& segment_caches) {
944
25
    DCHECK_EQ(block->columns(), _tablet_schema.num_columns());
945
    // the process logic here is the same as MemTable::_aggregate_for_flexible_partial_update_without_seq_col()
946
    // after this function, there will be at most 2 rows for a specified key
947
25
    std::vector<BitmapValue>* skip_bitmaps = &(
948
25
            assert_cast<ColumnBitmap*>(block->get_by_position(_tablet_schema.skip_bitmap_col_idx())
949
25
                                               .column->assume_mutable()
950
25
                                               .get())
951
25
                    ->get_data());
952
25
    const auto* delete_signs = BaseTablet::get_delete_sign_column_data(*block, num_rows);
953
954
25
    auto filtered_block = _tablet_schema.create_block();
955
25
    MutableBlock output_block = MutableBlock::build_mutable_block(&filtered_block);
956
957
25
    int same_key_rows {0};
958
25
    std::string previous_key {};
959
541
    for (int block_pos {0}; block_pos < num_rows; block_pos++) {
960
516
        std::string key = _writer._full_encode_keys(key_columns, block_pos);
961
516
        if (block_pos > 0 && previous_key == key) {
962
383
            same_key_rows++;
963
383
        } else {
964
133
            if (same_key_rows > 0) {
965
108
                RETURN_IF_ERROR(aggregate_rows(output_block, block, block_pos - same_key_rows,
966
108
                                               block_pos, std::move(previous_key), skip_bitmaps,
967
108
                                               delete_signs, seq_column, specified_rowsets,
968
108
                                               segment_caches));
969
108
            }
970
133
            same_key_rows = 1;
971
133
        }
972
516
        previous_key = std::move(key);
973
516
    }
974
25
    if (same_key_rows > 0) {
975
25
        RETURN_IF_ERROR(aggregate_rows(output_block, block, num_rows - same_key_rows, num_rows,
976
25
                                       std::move(previous_key), skip_bitmaps, delete_signs,
977
25
                                       seq_column, specified_rowsets, segment_caches));
978
25
    }
979
980
25
    block->swap(output_block.to_block());
981
25
    return Status::OK();
982
25
}
983
984
Status BlockAggregator::fill_sequence_column(Block* block, size_t num_rows,
985
                                             const FixedReadPlan& read_plan,
986
4
                                             std::vector<BitmapValue>& skip_bitmaps) {
987
4
    DCHECK(_tablet_schema.has_sequence_col());
988
4
    std::vector<uint32_t> cids {static_cast<uint32_t>(_tablet_schema.sequence_col_idx())};
989
4
    auto seq_col_unique_id = _tablet_schema.column(_tablet_schema.sequence_col_idx()).unique_id();
990
991
4
    auto seq_col_block = _tablet_schema.create_block_by_cids(cids);
992
4
    auto tmp_block = _tablet_schema.create_block_by_cids(cids);
993
4
    std::map<uint32_t, uint32_t> read_index;
994
4
    RETURN_IF_ERROR(read_plan.read_columns_by_plan(_tablet_schema, cids, _writer._rsid_to_rowset,
995
4
                                                   seq_col_block, &read_index, false));
996
997
4
    auto new_seq_col_ptr = tmp_block.get_by_position(0).column->assume_mutable();
998
4
    const auto& old_seq_col_ptr = *seq_col_block.get_by_position(0).column;
999
4
    const auto& cur_seq_col_ptr = *block->get_by_position(_tablet_schema.sequence_col_idx()).column;
1000
80
    for (uint32_t block_pos {0}; block_pos < num_rows; block_pos++) {
1001
76
        if (read_index.contains(block_pos)) {
1002
14
            new_seq_col_ptr->insert_from(old_seq_col_ptr, read_index[block_pos]);
1003
14
            skip_bitmaps[block_pos].remove(seq_col_unique_id);
1004
62
        } else {
1005
62
            new_seq_col_ptr->insert_from(cur_seq_col_ptr, block_pos);
1006
62
        }
1007
76
    }
1008
4
    block->replace_by_position(_tablet_schema.sequence_col_idx(), std::move(new_seq_col_ptr));
1009
4
    return Status::OK();
1010
4
}
1011
1012
Status BlockAggregator::aggregate_for_insert_after_delete(
1013
        Block* block, size_t num_rows, const std::vector<IOlapColumnDataAccessor*>& key_columns,
1014
        const std::vector<RowsetSharedPtr>& specified_rowsets,
1015
265
        std::vector<std::unique_ptr<SegmentCacheHandle>>& segment_caches) {
1016
265
    DCHECK_EQ(block->columns(), _tablet_schema.num_columns());
1017
    // there will be at most 2 rows for a specified key in block when control flow reaches here
1018
    // after this function, there will not be duplicate rows in block
1019
1020
265
    std::vector<BitmapValue>* skip_bitmaps = &(
1021
265
            assert_cast<ColumnBitmap*>(block->get_by_position(_tablet_schema.skip_bitmap_col_idx())
1022
265
                                               .column->assume_mutable()
1023
265
                                               .get())
1024
265
                    ->get_data());
1025
265
    const auto* delete_signs = BaseTablet::get_delete_sign_column_data(*block, num_rows);
1026
1027
265
    auto filter_column = ColumnUInt8::create(num_rows, 1);
1028
265
    auto* __restrict filter_map = filter_column->get_data().data();
1029
265
    std::string previous_key {};
1030
265
    bool previous_has_delete_sign {false};
1031
265
    int duplicate_rows {0};
1032
265
    int32_t delete_sign_col_unique_id =
1033
265
            _tablet_schema.column(_tablet_schema.delete_sign_idx()).unique_id();
1034
265
    auto seq_col_unique_id =
1035
265
            (_tablet_schema.sequence_col_idx() != -1)
1036
265
                    ? _tablet_schema.column(_tablet_schema.sequence_col_idx()).unique_id()
1037
265
                    : -1;
1038
265
    FixedReadPlan read_plan;
1039
2.45k
    for (size_t block_pos {0}; block_pos < num_rows; block_pos++) {
1040
2.19k
        size_t delta_pos = block_pos;
1041
2.19k
        auto& skip_bitmap = skip_bitmaps->at(block_pos);
1042
2.19k
        std::string key = _writer._full_encode_keys(key_columns, delta_pos);
1043
2.19k
        bool have_delete_sign =
1044
2.19k
                (!skip_bitmap.contains(delete_sign_col_unique_id) && delete_signs[block_pos] != 0);
1045
2.19k
        if (delta_pos > 0 && previous_key == key) {
1046
            // !!ATTENTION!!: We can only remove the row with delete sign if there is a insert with the same key after this row.
1047
            // If there is only a row with delete sign, we should keep it and can't remove it from block, because
1048
            // compaction will not use the delete bitmap when reading data. So there may still be rows with delete sign
1049
            // in later process
1050
43
            DCHECK(previous_has_delete_sign);
1051
43
            DCHECK(!have_delete_sign);
1052
43
            ++duplicate_rows;
1053
43
            RowLocation loc;
1054
43
            RowsetSharedPtr rowset;
1055
43
            Status st = _writer._tablet->lookup_row_key(
1056
43
                    key, &_tablet_schema, false, specified_rowsets, &loc,
1057
43
                    _writer._mow_context->max_version, segment_caches, &rowset, true);
1058
43
            bool is_expected_st = (st.is<ErrorCode::KEY_NOT_FOUND>() || st.ok());
1059
43
            DCHECK(is_expected_st || st.is<ErrorCode::MEM_LIMIT_EXCEEDED>())
1060
0
                    << "[BlockAggregator::aggregate_for_insert_after_delete] unexpected error "
1061
0
                       "status while lookup_row_key:"
1062
0
                    << st;
1063
43
            if (!is_expected_st) {
1064
0
                return st;
1065
0
            }
1066
1067
43
            Slice previous_seq_slice {};
1068
43
            if (st.ok()) {
1069
38
                if (_tablet_schema.has_sequence_col()) {
1070
                    // if the insert row doesn't specify the sequence column, we need to
1071
                    // read the historical's sequence column value so that we don't need
1072
                    // to handle seqeunce column in append_block_with_flexible_content()
1073
                    // for this row
1074
34
                    bool row_has_sequence_col = (!skip_bitmap.contains(seq_col_unique_id));
1075
34
                    if (!row_has_sequence_col) {
1076
14
                        read_plan.prepare_to_read(loc, block_pos);
1077
14
                        _writer._rsid_to_rowset.emplace(rowset->rowset_id(), rowset);
1078
14
                    }
1079
34
                }
1080
                // delete the existing row
1081
38
                _writer._mow_context->delete_bitmap->add(
1082
38
                        {loc.rowset_id, loc.segment_id, DeleteBitmap::TEMP_VERSION_COMMON},
1083
38
                        loc.row_id);
1084
38
            }
1085
            // and remove the row with delete sign from the current block
1086
43
            filter_map[block_pos - 1] = 0;
1087
43
        }
1088
2.19k
        previous_has_delete_sign = have_delete_sign;
1089
2.19k
        previous_key = std::move(key);
1090
2.19k
    }
1091
265
    if (duplicate_rows > 0) {
1092
9
        if (!read_plan.empty()) {
1093
            // fill sequence column value for some rows
1094
4
            RETURN_IF_ERROR(fill_sequence_column(block, num_rows, read_plan, *skip_bitmaps));
1095
4
        }
1096
9
        RETURN_IF_ERROR(filter_block(block, num_rows, std::move(filter_column), duplicate_rows,
1097
9
                                     "__filter_insert_after_delete_col__"));
1098
9
    }
1099
265
    return Status::OK();
1100
265
}
1101
1102
Status BlockAggregator::filter_block(Block* block, size_t num_rows, MutableColumnPtr filter_column,
1103
9
                                     int duplicate_rows, std::string col_name) {
1104
9
    auto num_cols = block->columns();
1105
9
    block->insert({std::move(filter_column), std::make_shared<DataTypeUInt8>(), col_name});
1106
9
    RETURN_IF_ERROR(Block::filter_block(block, num_cols, num_cols));
1107
9
    DCHECK_EQ(num_cols, block->columns());
1108
9
    size_t merged_rows = num_rows - block->rows();
1109
9
    if (duplicate_rows != merged_rows) {
1110
0
        auto msg = fmt::format(
1111
0
                "filter_block_for_flexible_partial_update {}: duplicate_rows != merged_rows, "
1112
0
                "duplicate_keys={}, merged_rows={}, num_rows={}, mutable_block->rows()={}",
1113
0
                col_name, duplicate_rows, merged_rows, num_rows, block->rows());
1114
0
        DCHECK(false) << msg;
1115
0
        return Status::InternalError<false>(msg);
1116
0
    }
1117
9
    return Status::OK();
1118
9
}
1119
1120
Status BlockAggregator::convert_pk_columns(Block* block, size_t row_pos, size_t num_rows,
1121
545
                                           std::vector<IOlapColumnDataAccessor*>& key_columns) {
1122
545
    key_columns.clear();
1123
1.14k
    for (uint32_t cid {0}; cid < _tablet_schema.num_key_columns(); cid++) {
1124
604
        RETURN_IF_ERROR(_writer._olap_data_convertor->set_source_content_with_specifid_column(
1125
604
                block->get_by_position(cid), row_pos, num_rows, cid));
1126
604
        auto [status, column] = _writer._olap_data_convertor->convert_column_data(cid);
1127
604
        if (!status.ok()) {
1128
0
            return status;
1129
0
        }
1130
604
        key_columns.push_back(column);
1131
604
    }
1132
545
    return Status::OK();
1133
545
}
1134
1135
Status BlockAggregator::convert_seq_column(Block* block, size_t row_pos, size_t num_rows,
1136
545
                                           IOlapColumnDataAccessor*& seq_column) {
1137
545
    seq_column = nullptr;
1138
545
    if (_tablet_schema.has_sequence_col()) {
1139
65
        auto seq_col_idx = _tablet_schema.sequence_col_idx();
1140
65
        RETURN_IF_ERROR(_writer._olap_data_convertor->set_source_content_with_specifid_column(
1141
65
                block->get_by_position(seq_col_idx), row_pos, num_rows, seq_col_idx));
1142
65
        auto [status, column] = _writer._olap_data_convertor->convert_column_data(seq_col_idx);
1143
65
        if (!status.ok()) {
1144
0
            return status;
1145
0
        }
1146
65
        seq_column = column;
1147
65
    }
1148
545
    return Status::OK();
1149
545
};
1150
1151
Status BlockAggregator::aggregate_for_flexible_partial_update(
1152
        Block* block, size_t num_rows, const std::vector<RowsetSharedPtr>& specified_rowsets,
1153
265
        std::vector<std::unique_ptr<SegmentCacheHandle>>& segment_caches) {
1154
265
    std::vector<IOlapColumnDataAccessor*> key_columns {};
1155
265
    IOlapColumnDataAccessor* seq_column {nullptr};
1156
1157
265
    RETURN_IF_ERROR(convert_pk_columns(block, 0, num_rows, key_columns));
1158
265
    RETURN_IF_ERROR(convert_seq_column(block, 0, num_rows, seq_column));
1159
1160
    // 1. merge duplicate rows when table has sequence column
1161
    // When there are multiple rows with the same keys in memtable, some of them specify specify the sequence column,
1162
    // some of them don't. We can't do the de-duplication in memtable because we don't know the historical data. We must
1163
    // de-duplicate them here.
1164
265
    if (_tablet_schema.has_sequence_col()) {
1165
25
        RETURN_IF_ERROR(aggregate_for_sequence_column(block, static_cast<int>(num_rows),
1166
25
                                                      key_columns, seq_column, specified_rowsets,
1167
25
                                                      segment_caches));
1168
25
    }
1169
1170
    // 2. merge duplicate rows and handle insert after delete
1171
265
    if (block->rows() != num_rows) {
1172
15
        num_rows = block->rows();
1173
        // data in block has changed, should re-encode key columns, sequence column
1174
15
        _writer._olap_data_convertor->clear_source_content();
1175
15
        RETURN_IF_ERROR(convert_pk_columns(block, 0, num_rows, key_columns));
1176
15
        RETURN_IF_ERROR(convert_seq_column(block, 0, num_rows, seq_column));
1177
15
    }
1178
265
    RETURN_IF_ERROR(aggregate_for_insert_after_delete(block, num_rows, key_columns,
1179
265
                                                      specified_rowsets, segment_caches));
1180
265
    return Status::OK();
1181
265
}
1182
1183
} // namespace doris