Coverage Report

Created: 2026-08-06 15:00

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/storage/segment/column_writer.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/segment/column_writer.h"
19
20
#include <gen_cpp/segment_v2.pb.h>
21
22
#include <algorithm>
23
#include <cstring>
24
#include <filesystem>
25
#include <memory>
26
27
#include "common/config.h"
28
#include "common/logging.h"
29
#include "core/data_type/data_type_agg_state.h"
30
#include "core/data_type/data_type_factory.hpp"
31
#include "core/types.h"
32
#include "io/fs/file_writer.h"
33
#include "storage/index/bloom_filter/bloom_filter_index_writer.h"
34
#include "storage/index/inverted/inverted_index_writer.h"
35
#include "storage/index/ordinal_page_index.h"
36
#include "storage/index/zone_map/zone_map_index.h"
37
#include "storage/olap_common.h"
38
#include "storage/segment/encoding_info.h"
39
#include "storage/segment/options.h"
40
#include "storage/segment/page_builder.h"
41
#include "storage/segment/page_io.h"
42
#include "storage/segment/page_pointer.h"
43
#include "storage/segment/variant/variant_column_writer_impl.h"
44
#include "storage/tablet/tablet_schema.h"
45
#include "storage/types.h"
46
#include "util/block_compression.h"
47
#include "util/debug_points.h"
48
#include "util/faststring.h"
49
#include "util/rle_encoding.h"
50
#include "util/simd/bits.h"
51
52
namespace doris::segment_v2 {
53
54
class NullBitmapBuilder {
55
public:
56
618k
    NullBitmapBuilder() : _has_null(false), _bitmap_buf(512), _rle_encoder(&_bitmap_buf, 1) {}
57
58
    explicit NullBitmapBuilder(size_t reserve_bits)
59
            : _has_null(false),
60
              _bitmap_buf(BitmapSize(reserve_bits)),
61
0
              _rle_encoder(&_bitmap_buf, 1) {}
62
63
3.03M
    void reserve_for_write(size_t num_rows, size_t non_null_count) {
64
3.03M
        if (num_rows == 0) {
65
0
            return;
66
0
        }
67
3.03M
        if (non_null_count == 0 || (non_null_count == num_rows && !_has_null)) {
68
501k
            if (_bitmap_buf.capacity() < kSmallReserveBytes) {
69
100
                _bitmap_buf.reserve(kSmallReserveBytes);
70
100
            }
71
501k
            return;
72
501k
        }
73
2.53M
        size_t raw_bytes = BitmapSize(num_rows);
74
2.53M
        size_t run_est = std::min(num_rows, non_null_count * 2 + 1);
75
2.53M
        size_t run_bytes_est = run_est * kBytesPerRun + kReserveSlackBytes;
76
2.53M
        size_t raw_overhead = raw_bytes / 63 + 1;
77
2.53M
        size_t raw_est = raw_bytes + raw_overhead + kReserveSlackBytes;
78
2.53M
        size_t reserve_bytes = std::min(raw_est, run_bytes_est);
79
2.53M
        if (_bitmap_buf.capacity() < reserve_bytes) {
80
9.87k
            const size_t cap = _bitmap_buf.capacity();
81
9.87k
            const size_t grow = cap + cap / 2;
82
9.87k
            const size_t new_cap = std::max(reserve_bytes, grow);
83
9.87k
            _bitmap_buf.reserve(new_cap);
84
9.87k
        }
85
2.53M
    }
86
87
8.48M
    void add_run(bool value, size_t run) {
88
8.48M
        _has_null |= value;
89
8.48M
        _rle_encoder.Put(value, run);
90
8.48M
    }
91
92
    // Returns whether the building nullmap contains nullptr
93
618k
    bool has_null() const { return _has_null; }
94
95
164k
    Status finish(OwnedSlice* slice) {
96
164k
        _rle_encoder.Flush();
97
164k
        RETURN_IF_CATCH_EXCEPTION({ *slice = _bitmap_buf.build(); });
98
164k
        return Status::OK();
99
164k
    }
100
101
618k
    void reset() {
102
618k
        _has_null = false;
103
618k
        _rle_encoder.Clear();
104
618k
    }
105
106
71.5k
    uint64_t size() { return _bitmap_buf.size(); }
107
108
private:
109
    static constexpr size_t kSmallReserveBytes = 64;
110
    static constexpr size_t kReserveSlackBytes = 16;
111
    static constexpr size_t kBytesPerRun = 6;
112
113
    bool _has_null;
114
    faststring _bitmap_buf;
115
    RleEncoder<bool> _rle_encoder;
116
};
117
118
inline ScalarColumnWriter* get_null_writer(const ColumnWriterOptions& opts,
119
76.6k
                                           io::FileWriter* file_writer, uint32_t id) {
120
76.6k
    if (!opts.meta->is_nullable()) {
121
33.2k
        return nullptr;
122
33.2k
    }
123
124
43.4k
    FieldType null_type = FieldType::OLAP_FIELD_TYPE_TINYINT;
125
43.4k
    ColumnWriterOptions null_options;
126
43.4k
    null_options.meta = opts.meta->add_children_columns();
127
43.4k
    null_options.meta->set_column_id(id);
128
43.4k
    null_options.meta->set_unique_id(id);
129
43.4k
    null_options.meta->set_type(int(null_type));
130
43.4k
    null_options.meta->set_is_nullable(false);
131
43.4k
    null_options.meta->set_length(
132
43.4k
            cast_set<int32_t>(field_type_size(FieldType::OLAP_FIELD_TYPE_TINYINT)));
133
43.4k
    null_options.meta->set_compression(opts.meta->compression());
134
135
43.4k
    null_options.need_zone_map = false;
136
43.4k
    null_options.need_bloom_filter = false;
137
43.4k
    null_options.storage_format = opts.storage_format;
138
139
43.4k
    auto null_column_ptr = std::make_shared<TabletColumn>(
140
43.4k
            FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE, null_type, false,
141
43.4k
            null_options.meta->unique_id(), null_options.meta->length());
142
43.4k
    null_column_ptr->set_name("nullable");
143
43.4k
    null_column_ptr->set_index_length(-1); // no short key index
144
43.4k
    null_options.meta->set_encoding(
145
43.4k
            EncodingInfo::resolve_default_encoding(opts.storage_format, *null_column_ptr));
146
43.4k
    return new ScalarColumnWriter(null_options, std::move(null_column_ptr), file_writer);
147
76.6k
}
148
149
ColumnWriter::ColumnWriter(TabletColumnPtr column, bool is_nullable, ColumnMetaPB* meta)
150
1.21M
        : _column(std::move(column)), _is_nullable(is_nullable), _column_meta(meta) {
151
1.21M
    _data_type = DataTypeFactory::instance().create_data_type(*_column_meta);
152
1.21M
}
153
Status ColumnWriter::create_struct_writer(const ColumnWriterOptions& opts,
154
                                          const TabletColumn* column, io::FileWriter* file_writer,
155
2.91k
                                          std::unique_ptr<ColumnWriter>* writer) {
156
    // not support empty struct
157
2.91k
    DCHECK(column->get_subtype_count() >= 1);
158
2.91k
    std::vector<std::unique_ptr<ColumnWriter>> sub_column_writers;
159
2.91k
    sub_column_writers.reserve(column->get_subtype_count());
160
22.3k
    for (uint32_t i = 0; i < column->get_subtype_count(); i++) {
161
19.4k
        const TabletColumn& sub_column = column->get_sub_column(i);
162
19.4k
        RETURN_IF_ERROR(sub_column.check_valid());
163
164
        // create sub writer
165
19.4k
        ColumnWriterOptions column_options;
166
19.4k
        column_options.meta = opts.meta->mutable_children_columns(i);
167
19.4k
        column_options.need_zone_map = false;
168
19.4k
        column_options.need_bloom_filter = sub_column.is_bf_column();
169
19.4k
        column_options.storage_format = opts.storage_format;
170
19.4k
        std::unique_ptr<ColumnWriter> sub_column_writer;
171
19.4k
        RETURN_IF_ERROR(
172
19.4k
                ColumnWriter::create(column_options, &sub_column, file_writer, &sub_column_writer));
173
19.4k
        sub_column_writers.push_back(std::move(sub_column_writer));
174
19.4k
    }
175
176
2.91k
    ScalarColumnWriter* null_writer =
177
2.91k
            get_null_writer(opts, file_writer, column->get_subtype_count() + 1);
178
179
2.91k
    *writer = std::unique_ptr<ColumnWriter>(new StructColumnWriter(
180
2.91k
            opts, std::make_shared<TabletColumn>(*column), null_writer, sub_column_writers));
181
2.91k
    return Status::OK();
182
2.91k
}
183
184
Status ColumnWriter::create_array_writer(const ColumnWriterOptions& opts,
185
                                         const TabletColumn* column, io::FileWriter* file_writer,
186
46.8k
                                         std::unique_ptr<ColumnWriter>* writer) {
187
46.8k
    DCHECK(column->get_subtype_count() == 1);
188
46.8k
    const TabletColumn& item_column = column->get_sub_column(0);
189
46.8k
    RETURN_IF_ERROR(item_column.check_valid());
190
191
    // create item writer
192
46.8k
    ColumnWriterOptions item_options;
193
46.8k
    item_options.meta = opts.meta->mutable_children_columns(0);
194
46.8k
    item_options.need_zone_map = false;
195
46.8k
    item_options.need_bloom_filter = item_column.is_bf_column();
196
46.8k
    item_options.storage_format = opts.storage_format;
197
46.8k
    std::unique_ptr<ColumnWriter> item_writer;
198
46.8k
    RETURN_IF_ERROR(ColumnWriter::create(item_options, &item_column, file_writer, &item_writer));
199
200
    // create length writer
201
46.8k
    FieldType length_type = FieldType::OLAP_FIELD_TYPE_UNSIGNED_BIGINT;
202
203
46.8k
    ColumnWriterOptions length_options;
204
46.8k
    length_options.meta = opts.meta->add_children_columns();
205
46.8k
    length_options.meta->set_column_id(2);
206
46.8k
    length_options.meta->set_unique_id(2);
207
46.8k
    length_options.meta->set_type(int(length_type));
208
46.8k
    length_options.meta->set_is_nullable(false);
209
46.8k
    length_options.meta->set_length(
210
46.8k
            cast_set<int32_t>(field_type_size(FieldType::OLAP_FIELD_TYPE_UNSIGNED_BIGINT)));
211
46.8k
    length_options.meta->set_compression(opts.meta->compression());
212
213
46.8k
    length_options.need_zone_map = false;
214
46.8k
    length_options.need_bloom_filter = false;
215
46.8k
    length_options.storage_format = opts.storage_format;
216
217
46.8k
    auto length_column_ptr = std::make_shared<TabletColumn>(
218
46.8k
            FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE, length_type,
219
46.8k
            length_options.meta->is_nullable(), length_options.meta->unique_id(),
220
46.8k
            length_options.meta->length());
221
46.8k
    length_column_ptr->set_name("length");
222
46.8k
    length_column_ptr->set_index_length(-1); // no short key index
223
46.8k
    length_options.meta->set_encoding(
224
46.8k
            EncodingInfo::resolve_default_encoding(opts.storage_format, *length_column_ptr));
225
46.8k
    auto* length_writer =
226
46.8k
            new OffsetColumnWriter(length_options, std::move(length_column_ptr), file_writer);
227
228
46.8k
    ScalarColumnWriter* null_writer = get_null_writer(opts, file_writer, 3);
229
230
46.8k
    *writer = std::unique_ptr<ColumnWriter>(
231
46.8k
            new ArrayColumnWriter(opts, std::make_shared<TabletColumn>(*column), length_writer,
232
46.8k
                                  null_writer, std::move(item_writer)));
233
46.8k
    return Status::OK();
234
46.8k
}
235
236
Status ColumnWriter::create_map_writer(const ColumnWriterOptions& opts, const TabletColumn* column,
237
                                       io::FileWriter* file_writer,
238
26.8k
                                       std::unique_ptr<ColumnWriter>* writer) {
239
26.8k
    DCHECK(column->get_subtype_count() == 2);
240
26.8k
    if (column->get_subtype_count() < 2) {
241
0
        return Status::InternalError(
242
0
                "If you upgraded from version 1.2.*, please DROP the MAP columns and then "
243
0
                "ADD the MAP columns back.");
244
0
    }
245
    // create key & value writer
246
26.8k
    std::vector<std::unique_ptr<ColumnWriter>> inner_writer_list;
247
80.6k
    for (int i = 0; i < 2; ++i) {
248
53.7k
        const TabletColumn& item_column = column->get_sub_column(i);
249
53.7k
        RETURN_IF_ERROR(item_column.check_valid());
250
251
        // create item writer
252
53.7k
        ColumnWriterOptions item_options;
253
53.7k
        item_options.meta = opts.meta->mutable_children_columns(i);
254
53.7k
        item_options.need_zone_map = false;
255
53.7k
        item_options.need_bloom_filter = item_column.is_bf_column();
256
53.7k
        item_options.storage_format = opts.storage_format;
257
53.7k
        std::unique_ptr<ColumnWriter> item_writer;
258
53.7k
        RETURN_IF_ERROR(
259
53.7k
                ColumnWriter::create(item_options, &item_column, file_writer, &item_writer));
260
53.7k
        inner_writer_list.push_back(std::move(item_writer));
261
53.7k
    }
262
263
    // create offset writer
264
26.8k
    FieldType length_type = FieldType::OLAP_FIELD_TYPE_UNSIGNED_BIGINT;
265
266
    // Be Cautious: column unique id is used for column reader creation
267
26.8k
    ColumnWriterOptions length_options;
268
26.8k
    length_options.meta = opts.meta->add_children_columns();
269
26.8k
    length_options.meta->set_column_id(column->get_subtype_count() + 1);
270
26.8k
    length_options.meta->set_unique_id(column->get_subtype_count() + 1);
271
26.8k
    length_options.meta->set_type(int(length_type));
272
26.8k
    length_options.meta->set_is_nullable(false);
273
26.8k
    length_options.meta->set_length(
274
26.8k
            cast_set<int32_t>(field_type_size(FieldType::OLAP_FIELD_TYPE_UNSIGNED_BIGINT)));
275
26.8k
    length_options.meta->set_compression(opts.meta->compression());
276
277
26.8k
    length_options.need_zone_map = false;
278
26.8k
    length_options.need_bloom_filter = false;
279
26.8k
    length_options.storage_format = opts.storage_format;
280
281
26.8k
    auto length_column_ptr = std::make_shared<TabletColumn>(
282
26.8k
            FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE, length_type,
283
26.8k
            length_options.meta->is_nullable(), length_options.meta->unique_id(),
284
26.8k
            length_options.meta->length());
285
26.8k
    length_column_ptr->set_name("length");
286
26.8k
    length_column_ptr->set_index_length(-1); // no short key index
287
26.8k
    length_options.meta->set_encoding(
288
26.8k
            EncodingInfo::resolve_default_encoding(opts.storage_format, *length_column_ptr));
289
26.8k
    auto* length_writer =
290
26.8k
            new OffsetColumnWriter(length_options, std::move(length_column_ptr), file_writer);
291
292
26.8k
    ScalarColumnWriter* null_writer =
293
26.8k
            get_null_writer(opts, file_writer, column->get_subtype_count() + 2);
294
295
26.8k
    *writer = std::unique_ptr<ColumnWriter>(
296
26.8k
            new MapColumnWriter(opts, std::make_shared<TabletColumn>(*column), null_writer,
297
26.8k
                                length_writer, inner_writer_list));
298
299
26.8k
    return Status::OK();
300
26.8k
}
301
302
Status ColumnWriter::create_agg_state_writer(const ColumnWriterOptions& opts,
303
                                             const TabletColumn* column,
304
                                             io::FileWriter* file_writer,
305
1.47k
                                             std::unique_ptr<ColumnWriter>* writer) {
306
1.47k
    auto data_type = DataTypeFactory::instance().create_data_type(*column);
307
1.47k
    const auto* agg_state_type = assert_cast<const DataTypeAggState*>(data_type.get());
308
1.47k
    auto type = agg_state_type->get_serialized_type()->get_primitive_type();
309
1.47k
    if (type == PrimitiveType::TYPE_STRING || type == PrimitiveType::INVALID_TYPE ||
310
1.47k
        type == PrimitiveType::TYPE_FIXED_LENGTH_OBJECT || type == PrimitiveType::TYPE_BITMAP) {
311
1.34k
        *writer = std::unique_ptr<ColumnWriter>(
312
1.34k
                new ScalarColumnWriter(opts, std::make_shared<TabletColumn>(*column), file_writer));
313
1.34k
    } else if (type == PrimitiveType::TYPE_ARRAY) {
314
53
        RETURN_IF_ERROR(create_array_writer(opts, column, file_writer, writer));
315
85
    } else if (type == PrimitiveType::TYPE_MAP) {
316
85
        RETURN_IF_ERROR(create_map_writer(opts, column, file_writer, writer));
317
18.4E
    } else {
318
18.4E
        throw Exception(ErrorCode::INTERNAL_ERROR,
319
18.4E
                        "OLAP_FIELD_TYPE_AGG_STATE meet unsupported type: {}",
320
18.4E
                        agg_state_type->get_name());
321
18.4E
    }
322
1.48k
    return Status::OK();
323
1.47k
}
324
325
Status ColumnWriter::create_variant_writer(const ColumnWriterOptions& opts,
326
                                           const TabletColumn* column, io::FileWriter* file_writer,
327
6.55k
                                           std::unique_ptr<ColumnWriter>* writer) {
328
    // Variant extracted columns have two kinds of physical writers:
329
    // - Doc-value snapshot column (`...__DORIS_VARIANT_DOC_VALUE__...`): use `VariantDocCompactWriter`
330
    //   to store the doc snapshot in a compact binary form.
331
    // - Regular extracted subcolumns: use `VariantSubcolumnWriter`.
332
    // The root VARIANT column itself uses `VariantColumnWriter`.
333
6.55k
    if (column->is_extracted_column()) {
334
490
        if (column->name().find(DOC_VALUE_COLUMN_PATH) != std::string::npos) {
335
382
            *writer = std::make_unique<VariantDocCompactWriter>(
336
382
                    opts, std::make_shared<TabletColumn>(*column));
337
382
            return Status::OK();
338
382
        }
339
108
        VLOG_DEBUG << "gen subwriter for " << column->path_info_ptr()->get_path();
340
108
        *writer = std::make_unique<VariantSubcolumnWriter>(opts,
341
108
                                                           std::make_shared<TabletColumn>(*column));
342
108
        return Status::OK();
343
490
    }
344
6.06k
    *writer = std::make_unique<VariantColumnWriter>(opts, std::make_shared<TabletColumn>(*column));
345
6.06k
    return Status::OK();
346
6.55k
}
347
348
//Todo(Amory): here should according nullable and offset and need sub to simply this function
349
Status ColumnWriter::create(const ColumnWriterOptions& opts, const TabletColumn* column,
350
1.07M
                            io::FileWriter* file_writer, std::unique_ptr<ColumnWriter>* writer) {
351
1.07M
    auto column_ptr = std::make_shared<TabletColumn>(*column);
352
1.07M
    if (is_scalar_type(column->type())) {
353
1.00M
        *writer = std::unique_ptr<ColumnWriter>(
354
1.00M
                new ScalarColumnWriter(opts, std::move(column_ptr), file_writer));
355
1.00M
        return Status::OK();
356
1.00M
    } else {
357
68.8k
        switch (column->type()) {
358
1.48k
        case FieldType::OLAP_FIELD_TYPE_AGG_STATE: {
359
1.48k
            RETURN_IF_ERROR(create_agg_state_writer(opts, column, file_writer, writer));
360
1.48k
            return Status::OK();
361
1.48k
        }
362
2.91k
        case FieldType::OLAP_FIELD_TYPE_STRUCT: {
363
2.91k
            RETURN_IF_ERROR(create_struct_writer(opts, column, file_writer, writer));
364
2.91k
            return Status::OK();
365
2.91k
        }
366
46.8k
        case FieldType::OLAP_FIELD_TYPE_ARRAY: {
367
46.8k
            RETURN_IF_ERROR(create_array_writer(opts, column, file_writer, writer));
368
46.8k
            return Status::OK();
369
46.8k
        }
370
12.1k
        case FieldType::OLAP_FIELD_TYPE_MAP: {
371
12.1k
            RETURN_IF_ERROR(create_map_writer(opts, column, file_writer, writer));
372
12.1k
            return Status::OK();
373
12.1k
        }
374
6.55k
        case FieldType::OLAP_FIELD_TYPE_VARIANT: {
375
            // Process columns with sparse column
376
6.55k
            RETURN_IF_ERROR(create_variant_writer(opts, column, file_writer, writer));
377
6.55k
            return Status::OK();
378
6.55k
        }
379
0
        default:
380
0
            return Status::NotSupported("unsupported type for ColumnWriter: {}",
381
0
                                        std::to_string(int(column_ptr->type())));
382
68.8k
        }
383
68.8k
    }
384
1.07M
}
385
386
Status ColumnWriter::append_nullable(const uint8_t* is_null_bits, const void* data,
387
700k
                                     size_t num_rows) {
388
700k
    const auto* ptr = (const uint8_t*)data;
389
700k
    BitmapIterator null_iter(is_null_bits, num_rows);
390
700k
    bool is_null = false;
391
700k
    size_t this_run = 0;
392
1.40M
    while ((this_run = null_iter.Next(&is_null)) > 0) {
393
700k
        if (is_null) {
394
1
            RETURN_IF_ERROR(append_nulls(this_run));
395
700k
        } else {
396
700k
            RETURN_IF_ERROR(append_data(&ptr, this_run));
397
700k
        }
398
700k
    }
399
700k
    return Status::OK();
400
700k
}
401
402
Status ColumnWriter::append_nullable(const uint8_t* null_map, const uint8_t** ptr,
403
0
                                     size_t num_rows) {
404
    // Fast path: use SIMD to detect all-NULL or all-non-NULL columns
405
0
    if (config::enable_rle_batch_put_optimization) {
406
0
        size_t non_null_count =
407
0
                simd::count_zero_num(reinterpret_cast<const int8_t*>(null_map), num_rows);
408
409
0
        if (non_null_count == 0) {
410
            // All NULL: skip run-length iteration, directly append all nulls
411
0
            RETURN_IF_ERROR(append_nulls(num_rows));
412
0
            *ptr += cell_size() * num_rows;
413
0
            return Status::OK();
414
0
        }
415
416
0
        if (non_null_count == num_rows) {
417
            // All non-NULL: skip run-length iteration, directly append all data
418
0
            return append_data(ptr, num_rows);
419
0
        }
420
0
    }
421
422
    // Mixed case or sparse optimization disabled: use run-length processing
423
0
    size_t offset = 0;
424
0
    auto next_run_step = [&]() {
425
0
        size_t step = 1;
426
0
        for (auto i = offset + 1; i < num_rows; ++i) {
427
0
            if (null_map[offset] == null_map[i]) {
428
0
                step++;
429
0
            } else {
430
0
                break;
431
0
            }
432
0
        }
433
0
        return step;
434
0
    };
435
436
0
    do {
437
0
        auto step = next_run_step();
438
0
        if (null_map[offset]) {
439
0
            RETURN_IF_ERROR(append_nulls(step));
440
0
            *ptr += cell_size() * step;
441
0
        } else {
442
            // TODO:
443
            //  1. `*ptr += cell_size() * step;` should do in this function, not append_data;
444
            //  2. support array vectorized load and ptr offset add
445
0
            RETURN_IF_ERROR(append_data(ptr, step));
446
0
        }
447
0
        offset += step;
448
0
    } while (offset < num_rows);
449
450
0
    return Status::OK();
451
0
}
452
453
3.52M
Status ColumnWriter::append(const uint8_t* nullmap, const void* data, size_t num_rows) {
454
3.52M
    assert(data && num_rows > 0);
455
3.52M
    const auto* ptr = (const uint8_t*)data;
456
3.52M
    if (nullmap) {
457
3.07M
        return append_nullable(nullmap, &ptr, num_rows);
458
3.07M
    }
459
447k
    if (is_nullable()) {
460
28
        _implicit_not_null_map.assign(num_rows, 0);
461
28
        return append_nullable(_implicit_not_null_map.data(), &ptr, num_rows);
462
28
    }
463
447k
    return append_data(&ptr, num_rows);
464
447k
}
465
466
///////////////////////////////////////////////////////////////////////////////////
467
468
ScalarColumnWriter::ScalarColumnWriter(const ColumnWriterOptions& opts, TabletColumnPtr column,
469
                                       io::FileWriter* file_writer)
470
1.13M
        : ColumnWriter(std::move(column), opts.meta->is_nullable(), opts.meta),
471
1.13M
          _opts(opts),
472
1.13M
          _file_writer(file_writer),
473
1.13M
          _data_size(0) {
474
    // these opts.meta fields should be set by client
475
1.13M
    DCHECK(opts.meta->has_column_id());
476
1.13M
    DCHECK(opts.meta->has_unique_id());
477
1.13M
    DCHECK(opts.meta->has_type());
478
1.13M
    DCHECK(opts.meta->has_length());
479
1.13M
    DCHECK(opts.meta->has_encoding());
480
1.13M
    DCHECK(opts.meta->has_compression());
481
1.13M
    DCHECK(opts.meta->has_is_nullable());
482
1.13M
    DCHECK(file_writer != nullptr);
483
1.13M
    _inverted_index_builders.resize(_opts.inverted_indexes.size());
484
1.13M
}
485
486
1.13M
ScalarColumnWriter::~ScalarColumnWriter() {
487
    // delete all pages
488
1.13M
    _pages.clear();
489
1.13M
}
490
491
1.12M
Status ScalarColumnWriter::init() {
492
1.12M
    RETURN_IF_ERROR(get_block_compression_codec(_opts.meta->compression(), &_compress_codec));
493
494
1.12M
    PageBuilder* page_builder = nullptr;
495
496
    // Caller must set a concrete (non-DEFAULT) encoding on the meta before init.
497
1.12M
    if (_opts.meta->encoding() == DEFAULT_ENCODING) {
498
0
        return Status::InternalError(
499
0
                "ColumnMetaPB encoding is DEFAULT_ENCODING for column_id={}, type={}; caller must "
500
0
                "resolve to a concrete encoding before ScalarColumnWriter::init",
501
0
                _opts.meta->column_id(), get_column()->type());
502
0
    }
503
1.12M
    RETURN_IF_ERROR(
504
1.12M
            EncodingInfo::get(get_column()->type(), _opts.meta->encoding(), &_encoding_info));
505
    // create page builder
506
1.12M
    PageBuilderOptions opts;
507
1.12M
    opts.data_page_size = _opts.data_page_size;
508
1.12M
    opts.dict_page_size = _opts.dict_page_size;
509
    // V3 segments store the dictionary word page (and the dict-overflow fallback plain page)
510
    // with the V3 binary plain layout; pre-V3 segments keep V1.
511
1.12M
    opts.dict_binary_plain_encoding =
512
1.12M
            (_opts.storage_format == TabletStorageFormatPB::TABLET_STORAGE_FORMAT_V3)
513
1.12M
                    ? PLAIN_ENCODING_V3
514
1.12M
                    : PLAIN_ENCODING;
515
1.12M
    RETURN_IF_ERROR(_encoding_info->create_page_builder(opts, &page_builder));
516
1.12M
    if (page_builder == nullptr) {
517
0
        return Status::NotSupported("Failed to create page builder for type {} and encoding {}",
518
0
                                    get_column()->type(), _opts.meta->encoding());
519
0
    }
520
18.4E
    VLOG_DEBUG << fmt::format(
521
18.4E
            "[verbose] scalar column writer init, column_id={}, type={}, encoding={}, "
522
18.4E
            "is_nullable={}",
523
18.4E
            _opts.meta->column_id(), get_column()->type(),
524
18.4E
            EncodingTypePB_Name(_opts.meta->encoding()), _opts.meta->is_nullable());
525
1.12M
    _page_builder.reset(page_builder);
526
    // create ordinal builder
527
1.12M
    _ordinal_index_builder = std::make_unique<OrdinalIndexWriter>();
528
    // create null bitmap builder
529
1.12M
    if (is_nullable()) {
530
618k
        _null_bitmap_builder = std::make_unique<NullBitmapBuilder>();
531
618k
    }
532
1.12M
    if (_opts.need_zone_map) {
533
838k
        RETURN_IF_ERROR(
534
838k
                ZoneMapIndexWriter::create(_data_type, get_column(), _zone_map_index_builder));
535
838k
    }
536
537
1.12M
    if (_opts.need_inverted_index) {
538
42.0k
        do {
539
84.3k
            for (size_t i = 0; i < _opts.inverted_indexes.size(); i++) {
540
42.2k
                DBUG_EXECUTE_IF("column_writer.init", {
541
42.2k
                    class InvertedIndexColumnWriterEmpty final : public IndexColumnWriter {
542
42.2k
                    public:
543
42.2k
                        Status init() override { return Status::OK(); }
544
42.2k
                        Status add_values(const std::string name, const void* values,
545
42.2k
                                          size_t count) override {
546
42.2k
                            return Status::OK();
547
42.2k
                        }
548
42.2k
                        Status add_array_values(size_t field_size, const void* value_ptr,
549
42.2k
                                                const uint8_t* null_map, const uint8_t* offsets_ptr,
550
42.2k
                                                size_t count) override {
551
42.2k
                            return Status::OK();
552
42.2k
                        }
553
42.2k
                        Status add_nulls(uint32_t count) override { return Status::OK(); }
554
42.2k
                        Status add_array_nulls(const uint8_t* null_map, size_t num_rows) override {
555
42.2k
                            return Status::OK();
556
42.2k
                        }
557
42.2k
                        Status finish() override { return Status::OK(); }
558
42.2k
                        int64_t size() const override { return 0; }
559
42.2k
                        void close_on_error() override {}
560
42.2k
                    };
561
562
42.2k
                    _inverted_index_builders[i] =
563
42.2k
                            std::make_unique<InvertedIndexColumnWriterEmpty>();
564
565
42.2k
                    break;
566
42.2k
                });
567
568
42.2k
                RETURN_IF_ERROR(IndexColumnWriter::create(
569
42.2k
                        get_column(), &_inverted_index_builders[i], _opts.index_file_writer,
570
42.2k
                        _opts.inverted_indexes[i]));
571
42.2k
            }
572
42.0k
        } while (false);
573
42.0k
    }
574
1.12M
    if (_opts.need_bloom_filter) {
575
8.38k
        if (_opts.is_ngram_bf_index) {
576
3.34k
            RETURN_IF_ERROR(NGramBloomFilterIndexWriterImpl::create(
577
3.34k
                    BloomFilterOptions(), get_column()->type(), _opts.gram_size, _opts.gram_bf_size,
578
3.34k
                    &_bloom_filter_index_builder));
579
5.03k
        } else {
580
5.03k
            RETURN_IF_ERROR(BloomFilterIndexWriter::create(_opts.bf_options, get_column()->type(),
581
5.03k
                                                           &_bloom_filter_index_builder));
582
5.03k
        }
583
8.38k
    }
584
1.12M
    return Status::OK();
585
1.12M
}
586
587
3.74M
Status ScalarColumnWriter::append_nulls(size_t num_rows) {
588
3.74M
    _null_bitmap_builder->add_run(true, num_rows);
589
3.74M
    _next_rowid += num_rows;
590
3.74M
    if (_opts.need_zone_map) {
591
3.66M
        _zone_map_index_builder->add_nulls(cast_set<uint32_t>(num_rows));
592
3.66M
    }
593
3.74M
    if (_opts.need_inverted_index) {
594
2.03M
        for (const auto& builder : _inverted_index_builders) {
595
2.03M
            RETURN_IF_ERROR(builder->add_nulls(cast_set<uint32_t>(num_rows)));
596
2.03M
        }
597
2.03M
    }
598
3.74M
    if (_opts.need_bloom_filter) {
599
759
        _bloom_filter_index_builder->add_nulls(cast_set<uint32_t>(num_rows));
600
759
    }
601
3.74M
    return Status::OK();
602
3.74M
}
603
604
// append data to page builder. this function will make sure that
605
// num_rows must be written before return. And ptr will be modified
606
// to next data should be written
607
5.24M
Status ScalarColumnWriter::append_data(const uint8_t** ptr, size_t num_rows) {
608
5.24M
    size_t remaining = num_rows;
609
10.5M
    while (remaining > 0) {
610
5.32M
        size_t num_written = remaining;
611
5.32M
        RETURN_IF_ERROR(append_data_in_current_page(ptr, &num_written));
612
613
5.32M
        remaining -= num_written;
614
615
5.32M
        if (_page_builder->is_page_full()) {
616
78.1k
            RETURN_IF_ERROR(finish_current_page());
617
78.1k
        }
618
5.32M
    }
619
5.24M
    return Status::OK();
620
5.24M
}
621
622
Status ScalarColumnWriter::_internal_append_data_in_current_page(const uint8_t* data,
623
5.39M
                                                                 size_t* num_written) {
624
5.39M
    RETURN_IF_ERROR(_page_builder->add(data, num_written));
625
5.39M
    if (_opts.need_zone_map) {
626
5.05M
        _zone_map_index_builder->add_values(data, *num_written);
627
5.05M
    }
628
5.39M
    if (_opts.need_inverted_index) {
629
2.05M
        for (const auto& builder : _inverted_index_builders) {
630
2.05M
            RETURN_IF_ERROR(builder->add_values(get_column()->name(), data, *num_written));
631
2.05M
        }
632
2.05M
    }
633
5.39M
    if (_opts.need_bloom_filter) {
634
8.96k
        RETURN_IF_ERROR(_bloom_filter_index_builder->add_values(data, *num_written));
635
8.96k
    }
636
637
5.39M
    _next_rowid += *num_written;
638
639
    // we must write null bits after write data, because we don't
640
    // know how many rows can be written into current page
641
5.39M
    if (is_nullable()) {
642
4.75M
        _null_bitmap_builder->add_run(false, *num_written);
643
4.75M
    }
644
5.39M
    return Status::OK();
645
5.39M
}
646
647
5.39M
Status ScalarColumnWriter::append_data_in_current_page(const uint8_t** data, size_t* num_written) {
648
5.39M
    RETURN_IF_ERROR(append_data_in_current_page(*data, num_written));
649
5.39M
    *data += cell_size() * (*num_written);
650
5.39M
    return Status::OK();
651
5.39M
}
652
653
Status ScalarColumnWriter::append_nullable(const uint8_t* null_map, const uint8_t** ptr,
654
3.03M
                                           size_t num_rows) {
655
    // When optimization is disabled, use base class implementation
656
3.03M
    if (!config::enable_rle_batch_put_optimization) {
657
0
        return ColumnWriter::append_nullable(null_map, ptr, num_rows);
658
0
    }
659
660
3.03M
    if (UNLIKELY(num_rows == 0)) {
661
0
        return Status::OK();
662
0
    }
663
664
    // Build run-length encoded null runs using memchr for fast boundary detection
665
3.03M
    _null_run_buffer.clear();
666
3.03M
    if (_null_run_buffer.capacity() < num_rows) {
667
630k
        _null_run_buffer.reserve(std::min(num_rows, size_t(256)));
668
630k
    }
669
670
3.03M
    size_t non_null_count = 0;
671
3.03M
    size_t offset = 0;
672
8.36M
    while (offset < num_rows) {
673
5.33M
        bool is_null = null_map[offset] != 0;
674
5.33M
        size_t remaining = num_rows - offset;
675
5.33M
        const uint8_t* run_end =
676
5.33M
                static_cast<const uint8_t*>(memchr(null_map + offset, is_null ? 0 : 1, remaining));
677
5.33M
        size_t run_length = run_end != nullptr ? (run_end - (null_map + offset)) : remaining;
678
5.33M
        _null_run_buffer.push_back(NullRun {is_null, static_cast<uint32_t>(run_length)});
679
5.33M
        if (!is_null) {
680
4.09M
            non_null_count += run_length;
681
4.09M
        }
682
5.33M
        offset += run_length;
683
5.33M
    }
684
685
    // Pre-allocate buffer based on estimated size
686
3.04M
    if (_null_bitmap_builder != nullptr) {
687
3.04M
        size_t current_rows = _next_rowid - _first_rowid;
688
3.04M
        size_t expected_rows = current_rows + num_rows;
689
3.04M
        size_t est_non_null = non_null_count;
690
3.04M
        if (num_rows > 0 && expected_rows > num_rows) {
691
2.48M
            est_non_null = (non_null_count * expected_rows) / num_rows;
692
2.48M
        }
693
3.04M
        _null_bitmap_builder->reserve_for_write(expected_rows, est_non_null);
694
3.04M
    }
695
696
3.03M
    if (non_null_count == 0) {
697
        // All NULL: skip data writing, only update null bitmap and indexes
698
53.7k
        RETURN_IF_ERROR(append_nulls(num_rows));
699
53.7k
        *ptr += cell_size() * num_rows;
700
53.7k
        return Status::OK();
701
53.7k
    }
702
703
2.97M
    if (non_null_count == num_rows) {
704
        // All non-NULL: use normal append_data which handles both data and null bitmap
705
2.92M
        return append_data(ptr, num_rows);
706
2.92M
    }
707
708
    // Process by runs
709
2.35M
    for (const auto& run : _null_run_buffer) {
710
2.35M
        size_t run_length = run.len;
711
2.35M
        if (run.is_null) {
712
1.18M
            RETURN_IF_ERROR(append_nulls(run_length));
713
1.18M
            *ptr += cell_size() * run_length;
714
1.18M
        } else {
715
            // TODO:
716
            //  1. `*ptr += cell_size() * step;` should do in this function, not append_data;
717
            //  2. support array vectorized load and ptr offset add
718
1.16M
            RETURN_IF_ERROR(append_data(ptr, run_length));
719
1.16M
        }
720
2.35M
    }
721
722
58.5k
    return Status::OK();
723
58.5k
}
724
725
136k
uint64_t ScalarColumnWriter::estimate_buffer_size() {
726
136k
    uint64_t size = _data_size;
727
136k
    size += _page_builder->size();
728
136k
    if (is_nullable()) {
729
71.5k
        size += _null_bitmap_builder->size();
730
71.5k
    }
731
136k
    size += _ordinal_index_builder->size();
732
136k
    if (_opts.need_zone_map) {
733
119k
        size += _zone_map_index_builder->size();
734
119k
    }
735
136k
    if (_opts.need_bloom_filter) {
736
267
        size += _bloom_filter_index_builder->size();
737
267
    }
738
136k
    return size;
739
136k
}
740
741
1.13M
Status ScalarColumnWriter::finish() {
742
1.13M
    RETURN_IF_ERROR(finish_current_page());
743
1.13M
    _opts.meta->set_num_rows(_next_rowid);
744
1.13M
    return Status::OK();
745
1.13M
}
746
747
1.13M
Status ScalarColumnWriter::write_data() {
748
1.13M
    auto offset = _file_writer->bytes_appended();
749
1.53M
    auto collect_uncompressed_bytes = [](const PageFooterPB& footer) {
750
1.53M
        return footer.uncompressed_size() + footer.ByteSizeLong() +
751
1.53M
               sizeof(uint32_t) /* footer size */ + sizeof(uint32_t) /* checksum */;
752
1.53M
    };
753
1.16M
    for (auto& page : _pages) {
754
1.16M
        _total_uncompressed_data_pages_size += collect_uncompressed_bytes(page->footer);
755
1.16M
        RETURN_IF_ERROR(_write_data_page(page.get()));
756
1.16M
    }
757
1.13M
    _pages.clear();
758
    // write column dict
759
1.13M
    if (_encoding_info->encoding() == DICT_ENCODING) {
760
366k
        OwnedSlice dict_body;
761
366k
        RETURN_IF_ERROR(_page_builder->get_dictionary_page(&dict_body));
762
366k
        EncodingTypePB dict_word_page_encoding;
763
366k
        RETURN_IF_ERROR(_page_builder->get_dictionary_page_encoding(&dict_word_page_encoding));
764
765
366k
        PageFooterPB footer;
766
366k
        footer.set_type(DICTIONARY_PAGE);
767
366k
        footer.set_uncompressed_size(cast_set<uint32_t>(dict_body.slice().get_size()));
768
366k
        footer.mutable_dict_page_footer()->set_encoding(dict_word_page_encoding);
769
366k
        _total_uncompressed_data_pages_size += collect_uncompressed_bytes(footer);
770
771
366k
        PagePointer dict_pp;
772
366k
        RETURN_IF_ERROR(PageIO::compress_and_write_page(
773
366k
                _compress_codec, _opts.compression_min_space_saving, _file_writer,
774
366k
                {dict_body.slice()}, footer, &dict_pp));
775
366k
        dict_pp.to_proto(_opts.meta->mutable_dict_page());
776
366k
    }
777
1.13M
    _total_compressed_data_pages_size += _file_writer->bytes_appended() - offset;
778
1.13M
    _page_builder.reset();
779
1.13M
    return Status::OK();
780
1.13M
}
781
782
1.08M
Status ScalarColumnWriter::write_ordinal_index() {
783
1.08M
    return _ordinal_index_builder->finish(_file_writer, _opts.meta->add_indexes());
784
1.08M
}
785
786
891k
Status ScalarColumnWriter::write_zone_map() {
787
891k
    if (_opts.need_zone_map) {
788
838k
        return _zone_map_index_builder->finish(_file_writer, _opts.meta->add_indexes());
789
838k
    }
790
53.6k
    return Status::OK();
791
891k
}
792
793
790k
Status ScalarColumnWriter::write_inverted_index() {
794
790k
    if (_opts.need_inverted_index) {
795
42.2k
        for (const auto& builder : _inverted_index_builders) {
796
42.2k
            RETURN_IF_ERROR(builder->finish());
797
42.2k
        }
798
42.0k
    }
799
790k
    return Status::OK();
800
790k
}
801
802
782k
Status ScalarColumnWriter::write_bloom_filter_index() {
803
782k
    if (_opts.need_bloom_filter) {
804
8.38k
        return _bloom_filter_index_builder->finish(_file_writer, _opts.meta->add_indexes());
805
8.38k
    }
806
773k
    return Status::OK();
807
782k
}
808
809
// write a data page into file and update ordinal index
810
1.16M
Status ScalarColumnWriter::_write_data_page(Page* page) {
811
1.16M
    PagePointer pp;
812
1.16M
    std::vector<Slice> compressed_body;
813
2.22M
    for (auto& data : page->data) {
814
2.22M
        compressed_body.push_back(data.slice());
815
2.22M
    }
816
1.16M
    RETURN_IF_ERROR(PageIO::write_page(_file_writer, compressed_body, page->footer, &pp));
817
1.16M
    _ordinal_index_builder->append_entry(page->footer.data_page_footer().first_ordinal(), pp);
818
1.16M
    return Status::OK();
819
1.16M
}
820
821
1.21M
Status ScalarColumnWriter::finish_current_page() {
822
1.21M
    if (_next_rowid == _first_rowid) {
823
44.1k
        return Status::OK();
824
44.1k
    }
825
1.16M
    if (_opts.need_zone_map) {
826
        // If the number of rows in the current page is less than the threshold,
827
        // we will invalidate zone map index for this page by set pass_all to true.
828
906k
        if (_next_rowid - _first_rowid < config::zone_map_row_num_threshold) {
829
643k
            _zone_map_index_builder->invalid_page_zone_map();
830
643k
        }
831
906k
        RETURN_IF_ERROR(_zone_map_index_builder->flush());
832
906k
    }
833
834
1.16M
    if (_opts.need_bloom_filter) {
835
8.75k
        RETURN_IF_ERROR(_bloom_filter_index_builder->flush());
836
8.75k
    }
837
838
1.16M
    _raw_data_bytes += _page_builder->get_raw_data_size();
839
840
    // build data page body : encoded values + [nullmap]
841
1.16M
    std::vector<Slice> body;
842
1.16M
    OwnedSlice encoded_values;
843
1.16M
    RETURN_IF_ERROR(_page_builder->finish(&encoded_values));
844
1.16M
    RETURN_IF_ERROR(_page_builder->reset());
845
1.16M
    body.push_back(encoded_values.slice());
846
847
1.16M
    OwnedSlice nullmap;
848
1.16M
    if (_null_bitmap_builder != nullptr) {
849
618k
        if (is_nullable() && _null_bitmap_builder->has_null()) {
850
164k
            RETURN_IF_ERROR(_null_bitmap_builder->finish(&nullmap));
851
164k
            body.push_back(nullmap.slice());
852
164k
        }
853
618k
        _null_bitmap_builder->reset();
854
618k
    }
855
856
    // prepare data page footer
857
1.16M
    std::unique_ptr<Page> page(new Page());
858
1.16M
    page->footer.set_type(DATA_PAGE);
859
1.16M
    page->footer.set_uncompressed_size(cast_set<uint32_t>(Slice::compute_total_size(body)));
860
1.16M
    auto* data_page_footer = page->footer.mutable_data_page_footer();
861
1.16M
    data_page_footer->set_first_ordinal(_first_rowid);
862
1.16M
    data_page_footer->set_num_values(_next_rowid - _first_rowid);
863
1.16M
    data_page_footer->set_nullmap_size(cast_set<uint32_t>(nullmap.slice().size));
864
1.16M
    if (_new_page_callback != nullptr) {
865
73.6k
        _new_page_callback->put_extra_info_in_page(data_page_footer);
866
73.6k
    }
867
    // trying to compress page body
868
1.16M
    OwnedSlice compressed_body;
869
1.16M
    RETURN_IF_ERROR(PageIO::compress_page_body(_compress_codec, _opts.compression_min_space_saving,
870
1.16M
                                               body, &compressed_body));
871
1.16M
    if (compressed_body.slice().empty()) {
872
        // page body is uncompressed
873
1.05M
        page->data.emplace_back(std::move(encoded_values));
874
1.05M
        page->data.emplace_back(std::move(nullmap));
875
1.05M
    } else {
876
        // page body is compressed
877
113k
        page->data.emplace_back(std::move(compressed_body));
878
113k
    }
879
880
1.16M
    _push_back_page(std::move(page));
881
1.16M
    _first_rowid = _next_rowid;
882
1.16M
    return Status::OK();
883
1.16M
}
884
885
////////////////////////////////////////////////////////////////////////////////
886
887
////////////////////////////////////////////////////////////////////////////////
888
// offset column writer
889
////////////////////////////////////////////////////////////////////////////////
890
891
OffsetColumnWriter::OffsetColumnWriter(const ColumnWriterOptions& opts, TabletColumnPtr column,
892
                                       io::FileWriter* file_writer)
893
73.7k
        : ScalarColumnWriter(opts, std::move(column), file_writer) {
894
    // now we only explain data in offset column as uint64
895
73.7k
    DCHECK(get_column()->type() == FieldType::OLAP_FIELD_TYPE_UNSIGNED_BIGINT);
896
73.7k
}
897
898
73.8k
OffsetColumnWriter::~OffsetColumnWriter() = default;
899
900
73.7k
Status OffsetColumnWriter::init() {
901
73.7k
    RETURN_IF_ERROR(ScalarColumnWriter::init());
902
73.7k
    register_flush_page_callback(this);
903
73.7k
    _next_offset = 0;
904
73.7k
    return Status::OK();
905
73.7k
}
906
907
73.8k
Status OffsetColumnWriter::append_data(const uint8_t** ptr, size_t num_rows) {
908
73.8k
    size_t remaining = num_rows;
909
148k
    while (remaining > 0) {
910
74.3k
        size_t num_written = remaining;
911
74.3k
        RETURN_IF_ERROR(append_data_in_current_page(ptr, &num_written));
912
        // Callers provide one extra tail offset after the written rows so the page footer can
913
        // store the next array item ordinal for the current page.
914
74.3k
        _next_offset = *(const uint64_t*)(*ptr);
915
74.3k
        remaining -= num_written;
916
917
74.3k
        if (_page_builder->is_page_full()) {
918
            // get next data for next array_item_rowid
919
606
            RETURN_IF_ERROR(finish_current_page());
920
606
        }
921
74.3k
    }
922
73.8k
    return Status::OK();
923
73.8k
}
924
925
73.6k
void OffsetColumnWriter::put_extra_info_in_page(DataPageFooterPB* footer) {
926
73.6k
    footer->set_next_array_item_ordinal(_next_offset);
927
73.6k
}
928
929
StructColumnWriter::StructColumnWriter(
930
        const ColumnWriterOptions& opts, TabletColumnPtr column, ScalarColumnWriter* null_writer,
931
        std::vector<std::unique_ptr<ColumnWriter>>& sub_column_writers)
932
2.91k
        : ColumnWriter(std::move(column), opts.meta->is_nullable(), opts.meta), _opts(opts) {
933
19.4k
    for (auto& sub_column_writer : sub_column_writers) {
934
19.4k
        _sub_column_writers.push_back(std::move(sub_column_writer));
935
19.4k
    }
936
2.91k
    _num_sub_column_writers = _sub_column_writers.size();
937
2.91k
    DCHECK(_num_sub_column_writers >= 1);
938
2.91k
    if (is_nullable()) {
939
2.62k
        _null_writer.reset(null_writer);
940
2.62k
    }
941
2.91k
}
942
943
2.91k
Status StructColumnWriter::init() {
944
19.4k
    for (auto& column_writer : _sub_column_writers) {
945
19.4k
        RETURN_IF_ERROR(column_writer->init());
946
19.4k
    }
947
2.91k
    if (is_nullable()) {
948
2.62k
        RETURN_IF_ERROR(_null_writer->init());
949
2.62k
    }
950
2.91k
    return Status::OK();
951
2.91k
}
952
953
2.28k
Status StructColumnWriter::write_inverted_index() {
954
2.28k
    if (_opts.need_inverted_index) {
955
0
        for (auto& column_writer : _sub_column_writers) {
956
0
            RETURN_IF_ERROR(column_writer->write_inverted_index());
957
0
        }
958
0
    }
959
2.28k
    return Status::OK();
960
2.28k
}
961
962
Status StructColumnWriter::append_nullable(const uint8_t* null_map, const uint8_t** ptr,
963
2.58k
                                           size_t num_rows) {
964
2.58k
    RETURN_IF_ERROR(append_data(ptr, num_rows));
965
2.58k
    RETURN_IF_ERROR(_null_writer->append_data(&null_map, num_rows));
966
2.58k
    return Status::OK();
967
2.58k
}
968
969
2.87k
Status StructColumnWriter::append_data(const uint8_t** ptr, size_t num_rows) {
970
2.87k
    const auto* results = reinterpret_cast<const uint64_t*>(*ptr);
971
21.9k
    for (size_t i = 0; i < _num_sub_column_writers; ++i) {
972
19.1k
        auto nullmap = *(results + _num_sub_column_writers + i);
973
19.1k
        auto data = *(results + i);
974
19.1k
        RETURN_IF_ERROR(_sub_column_writers[i]->append(reinterpret_cast<const uint8_t*>(nullmap),
975
19.1k
                                                       reinterpret_cast<const void*>(data),
976
19.1k
                                                       num_rows));
977
19.1k
    }
978
2.87k
    return Status::OK();
979
2.87k
}
980
981
188
uint64_t StructColumnWriter::estimate_buffer_size() {
982
188
    uint64_t size = 0;
983
784
    for (auto& column_writer : _sub_column_writers) {
984
784
        size += column_writer->estimate_buffer_size();
985
784
    }
986
188
    size += is_nullable() ? _null_writer->estimate_buffer_size() : 0;
987
188
    return size;
988
188
}
989
990
2.91k
Status StructColumnWriter::finish() {
991
19.4k
    for (auto& column_writer : _sub_column_writers) {
992
19.4k
        RETURN_IF_ERROR(column_writer->finish());
993
19.4k
    }
994
2.91k
    if (is_nullable()) {
995
2.62k
        RETURN_IF_ERROR(_null_writer->finish());
996
2.62k
    }
997
2.91k
    _opts.meta->set_num_rows(get_next_rowid());
998
2.91k
    return Status::OK();
999
2.91k
}
1000
1001
2.91k
Status StructColumnWriter::write_data() {
1002
19.4k
    for (auto& column_writer : _sub_column_writers) {
1003
19.4k
        RETURN_IF_ERROR(column_writer->write_data());
1004
19.4k
    }
1005
2.91k
    if (is_nullable()) {
1006
2.62k
        RETURN_IF_ERROR(_null_writer->write_data());
1007
2.62k
    }
1008
2.91k
    return Status::OK();
1009
2.91k
}
1010
1011
2.87k
Status StructColumnWriter::write_ordinal_index() {
1012
19.1k
    for (auto& column_writer : _sub_column_writers) {
1013
19.1k
        RETURN_IF_ERROR(column_writer->write_ordinal_index());
1014
19.1k
    }
1015
2.87k
    if (is_nullable()) {
1016
2.58k
        RETURN_IF_ERROR(_null_writer->write_ordinal_index());
1017
2.58k
    }
1018
2.87k
    return Status::OK();
1019
2.87k
}
1020
1021
0
Status StructColumnWriter::append_nulls(size_t num_rows) {
1022
0
    for (auto& column_writer : _sub_column_writers) {
1023
0
        RETURN_IF_ERROR(column_writer->append_nulls(num_rows));
1024
0
    }
1025
0
    if (is_nullable()) {
1026
0
        std::vector<UInt8> null_signs(num_rows, 1);
1027
0
        const uint8_t* null_sign_ptr = null_signs.data();
1028
0
        RETURN_IF_ERROR(_null_writer->append_data(&null_sign_ptr, num_rows));
1029
0
    }
1030
0
    return Status::OK();
1031
0
}
1032
1033
0
Status StructColumnWriter::finish_current_page() {
1034
0
    return Status::NotSupported("struct writer has no data, can not finish_current_page");
1035
0
}
1036
1037
ArrayColumnWriter::ArrayColumnWriter(const ColumnWriterOptions& opts, TabletColumnPtr column,
1038
                                     OffsetColumnWriter* offset_writer,
1039
                                     ScalarColumnWriter* null_writer,
1040
                                     std::unique_ptr<ColumnWriter> item_writer)
1041
46.8k
        : ColumnWriter(std::move(column), opts.meta->is_nullable(), opts.meta),
1042
46.8k
          _item_writer(std::move(item_writer)),
1043
46.8k
          _opts(opts) {
1044
46.8k
    _offset_writer.reset(offset_writer);
1045
46.8k
    if (is_nullable()) {
1046
31.0k
        _null_writer.reset(null_writer);
1047
31.0k
    }
1048
46.8k
}
1049
1050
46.8k
Status ArrayColumnWriter::init() {
1051
46.8k
    RETURN_IF_ERROR(_offset_writer->init());
1052
46.8k
    if (is_nullable()) {
1053
31.1k
        RETURN_IF_ERROR(_null_writer->init());
1054
31.1k
    }
1055
46.8k
    RETURN_IF_ERROR(_item_writer->init());
1056
46.8k
    if (_opts.need_inverted_index) {
1057
1.81k
        auto* writer = dynamic_cast<ScalarColumnWriter*>(_item_writer.get());
1058
1.81k
        if (writer != nullptr) {
1059
1.81k
            RETURN_IF_ERROR(IndexColumnWriter::create(get_column(), &_inverted_index_writer,
1060
1.81k
                                                      _opts.index_file_writer,
1061
1.81k
                                                      _opts.inverted_indexes[0]));
1062
1.81k
        }
1063
1.81k
    }
1064
46.8k
    if (_opts.need_ann_index) {
1065
76
        auto* writer = dynamic_cast<ScalarColumnWriter*>(_item_writer.get());
1066
76
        if (writer != nullptr) {
1067
76
            _ann_index_writer = std::make_unique<AnnIndexColumnWriter>(_opts.index_file_writer,
1068
76
                                                                       _opts.ann_index);
1069
76
            RETURN_IF_ERROR(_ann_index_writer->init());
1070
76
        }
1071
76
    }
1072
46.8k
    return Status::OK();
1073
46.8k
}
1074
1075
41.4k
Status ArrayColumnWriter::write_inverted_index() {
1076
41.4k
    if (_opts.need_inverted_index) {
1077
1.81k
        return _inverted_index_writer->finish();
1078
1.81k
    }
1079
39.6k
    return Status::OK();
1080
41.4k
}
1081
1082
41.3k
Status ArrayColumnWriter::write_ann_index() {
1083
41.3k
    if (_opts.need_ann_index) {
1084
73
        return _ann_index_writer->finish();
1085
73
    }
1086
41.3k
    return Status::OK();
1087
41.3k
}
1088
1089
// batch append data for array
1090
46.6k
Status ArrayColumnWriter::append_data(const uint8_t** ptr, size_t num_rows) {
1091
    // data_ptr contains
1092
    // [size, offset_ptr, item_data_ptr, item_nullmap_ptr]
1093
46.6k
    auto data_ptr = reinterpret_cast<const uint64_t*>(*ptr);
1094
    // total number length
1095
46.6k
    size_t element_cnt = size_t((unsigned long)(*data_ptr));
1096
46.6k
    auto offset_data = *(data_ptr + 1);
1097
46.6k
    const uint8_t* offsets_ptr = (const uint8_t*)offset_data;
1098
46.6k
    auto data = *(data_ptr + 2);
1099
46.6k
    auto nested_null_map = *(data_ptr + 3);
1100
46.6k
    if (element_cnt > 0) {
1101
30.0k
        RETURN_IF_ERROR(_item_writer->append(reinterpret_cast<const uint8_t*>(nested_null_map),
1102
30.0k
                                             reinterpret_cast<const void*>(data), element_cnt));
1103
30.0k
    }
1104
46.6k
    if (_opts.need_inverted_index) {
1105
1.81k
        auto* writer = dynamic_cast<ScalarColumnWriter*>(_item_writer.get());
1106
        // now only support nested type is scala
1107
1.81k
        if (writer != nullptr) {
1108
            //NOTE: use array field name as index field, but item_writer size should be used when moving item_data_ptr
1109
1.81k
            RETURN_IF_ERROR(_inverted_index_writer->add_array_values(
1110
1.81k
                    field_type_size(_item_writer->get_column()->type()),
1111
1.81k
                    reinterpret_cast<const void*>(data),
1112
1.81k
                    reinterpret_cast<const uint8_t*>(nested_null_map), offsets_ptr, num_rows));
1113
1.81k
        }
1114
1.81k
    }
1115
1116
46.6k
    if (_opts.need_ann_index) {
1117
76
        auto* writer = dynamic_cast<ScalarColumnWriter*>(_item_writer.get());
1118
        // now only support nested type is scala
1119
76
        if (writer != nullptr) {
1120
            //NOTE: use array field name as index field, but item_writer size should be used when moving item_data_ptr
1121
76
            RETURN_IF_ERROR(_ann_index_writer->add_array_values(
1122
76
                    field_type_size(_item_writer->get_column()->type()),
1123
76
                    reinterpret_cast<const void*>(data),
1124
76
                    reinterpret_cast<const uint8_t*>(nested_null_map), offsets_ptr, num_rows));
1125
76
        } else {
1126
0
            return Status::NotSupported(
1127
0
                    "Ann index can only be build on array with scalar type. but got {} as "
1128
0
                    "nested",
1129
0
                    _item_writer->get_column()->type());
1130
0
        }
1131
76
    }
1132
1133
46.6k
    RETURN_IF_ERROR(_offset_writer->append_data(&offsets_ptr, num_rows));
1134
46.6k
    return Status::OK();
1135
46.6k
}
1136
1137
1.14k
uint64_t ArrayColumnWriter::estimate_buffer_size() {
1138
1.14k
    return _offset_writer->estimate_buffer_size() +
1139
1.14k
           (is_nullable() ? _null_writer->estimate_buffer_size() : 0) +
1140
1.14k
           _item_writer->estimate_buffer_size();
1141
1.14k
}
1142
1143
Status ArrayColumnWriter::append_nullable(const uint8_t* null_map, const uint8_t** ptr,
1144
30.7k
                                          size_t num_rows) {
1145
30.7k
    RETURN_IF_ERROR(append_data(ptr, num_rows));
1146
30.7k
    if (is_nullable()) {
1147
30.7k
        if (_opts.need_inverted_index) {
1148
1.41k
            RETURN_IF_ERROR(_inverted_index_writer->add_array_nulls(null_map, num_rows));
1149
1.41k
        }
1150
30.7k
        RETURN_IF_ERROR(_null_writer->append_data(&null_map, num_rows));
1151
30.7k
    }
1152
30.7k
    return Status::OK();
1153
30.7k
}
1154
1155
46.9k
Status ArrayColumnWriter::finish() {
1156
46.9k
    RETURN_IF_ERROR(_offset_writer->finish());
1157
46.9k
    if (is_nullable()) {
1158
31.1k
        RETURN_IF_ERROR(_null_writer->finish());
1159
31.1k
    }
1160
46.9k
    RETURN_IF_ERROR(_item_writer->finish());
1161
46.9k
    _opts.meta->set_num_rows(get_next_rowid());
1162
46.9k
    return Status::OK();
1163
46.9k
}
1164
1165
46.9k
Status ArrayColumnWriter::write_data() {
1166
46.9k
    RETURN_IF_ERROR(_offset_writer->write_data());
1167
46.9k
    if (is_nullable()) {
1168
31.1k
        RETURN_IF_ERROR(_null_writer->write_data());
1169
31.1k
    }
1170
46.9k
    RETURN_IF_ERROR(_item_writer->write_data());
1171
46.9k
    return Status::OK();
1172
46.9k
}
1173
1174
46.4k
Status ArrayColumnWriter::write_ordinal_index() {
1175
46.4k
    RETURN_IF_ERROR(_offset_writer->write_ordinal_index());
1176
46.4k
    if (is_nullable()) {
1177
30.6k
        RETURN_IF_ERROR(_null_writer->write_ordinal_index());
1178
30.6k
    }
1179
46.4k
    if (!has_empty_items()) {
1180
29.9k
        RETURN_IF_ERROR(_item_writer->write_ordinal_index());
1181
29.9k
    }
1182
46.4k
    return Status::OK();
1183
46.4k
}
1184
1185
0
Status ArrayColumnWriter::append_nulls(size_t num_rows) {
1186
0
    const UInt64 offset = cast_set<UInt64>(_item_writer->get_next_rowid());
1187
0
    std::vector<UInt64> offsets_data(num_rows + 1, offset);
1188
0
    const uint8_t* offsets_ptr = reinterpret_cast<const uint8_t*>(offsets_data.data());
1189
0
    RETURN_IF_ERROR(_offset_writer->append_data(&offsets_ptr, num_rows));
1190
0
    return write_null_column(num_rows, true);
1191
0
}
1192
1193
0
Status ArrayColumnWriter::write_null_column(size_t num_rows, bool is_null) {
1194
0
    uint8_t null_sign = is_null ? 1 : 0;
1195
0
    while (is_nullable() && num_rows > 0) {
1196
        // TODO llj bulk write
1197
0
        const uint8_t* null_sign_ptr = &null_sign;
1198
0
        RETURN_IF_ERROR(_null_writer->append_data(&null_sign_ptr, 1));
1199
0
        --num_rows;
1200
0
    }
1201
0
    return Status::OK();
1202
0
}
1203
1204
0
Status ArrayColumnWriter::finish_current_page() {
1205
0
    return Status::NotSupported("array writer has no data, can not finish_current_page");
1206
0
}
1207
1208
/// ============================= MapColumnWriter =====================////
1209
MapColumnWriter::MapColumnWriter(const ColumnWriterOptions& opts, TabletColumnPtr column,
1210
                                 ScalarColumnWriter* null_writer, OffsetColumnWriter* offset_writer,
1211
                                 std::vector<std::unique_ptr<ColumnWriter>>& kv_writers)
1212
26.8k
        : ColumnWriter(std::move(column), opts.meta->is_nullable(), opts.meta), _opts(opts) {
1213
26.8k
    CHECK_EQ(kv_writers.size(), 2);
1214
26.8k
    _offsets_writer.reset(offset_writer);
1215
26.8k
    if (is_nullable()) {
1216
9.79k
        _null_writer.reset(null_writer);
1217
9.79k
    }
1218
53.7k
    for (auto& sub_writers : kv_writers) {
1219
53.7k
        _kv_writers.push_back(std::move(sub_writers));
1220
53.7k
    }
1221
26.8k
}
1222
1223
26.8k
Status MapColumnWriter::init() {
1224
26.8k
    RETURN_IF_ERROR(_offsets_writer->init());
1225
26.8k
    if (is_nullable()) {
1226
9.81k
        RETURN_IF_ERROR(_null_writer->init());
1227
9.81k
    }
1228
    // here register_flush_page_callback to call this.put_extra_info_in_page()
1229
    // when finish cur data page
1230
53.8k
    for (auto& sub_writer : _kv_writers) {
1231
53.8k
        RETURN_IF_ERROR(sub_writer->init());
1232
53.8k
    }
1233
26.8k
    return Status::OK();
1234
26.8k
}
1235
1236
1.58k
uint64_t MapColumnWriter::estimate_buffer_size() {
1237
1.58k
    size_t estimate = 0;
1238
3.16k
    for (auto& sub_writer : _kv_writers) {
1239
3.16k
        estimate += sub_writer->estimate_buffer_size();
1240
3.16k
    }
1241
1.58k
    estimate += _offsets_writer->estimate_buffer_size();
1242
1.58k
    if (is_nullable()) {
1243
1.55k
        estimate += _null_writer->estimate_buffer_size();
1244
1.55k
    }
1245
1.58k
    return estimate;
1246
1.58k
}
1247
1248
26.9k
Status MapColumnWriter::finish() {
1249
26.9k
    RETURN_IF_ERROR(_offsets_writer->finish());
1250
26.9k
    if (is_nullable()) {
1251
9.81k
        RETURN_IF_ERROR(_null_writer->finish());
1252
9.81k
    }
1253
53.7k
    for (auto& sub_writer : _kv_writers) {
1254
53.7k
        RETURN_IF_ERROR(sub_writer->finish());
1255
53.7k
    }
1256
26.9k
    _opts.meta->set_num_rows(get_next_rowid());
1257
26.9k
    return Status::OK();
1258
26.9k
}
1259
1260
Status MapColumnWriter::append_nullable(const uint8_t* null_map, const uint8_t** ptr,
1261
9.85k
                                        size_t num_rows) {
1262
9.85k
    RETURN_IF_ERROR(append_data(ptr, num_rows));
1263
9.86k
    if (is_nullable()) {
1264
9.86k
        RETURN_IF_ERROR(_null_writer->append_data(&null_map, num_rows));
1265
9.86k
    }
1266
9.85k
    return Status::OK();
1267
9.85k
}
1268
1269
// write key value data with offsets
1270
27.2k
Status MapColumnWriter::append_data(const uint8_t** ptr, size_t num_rows) {
1271
    // data_ptr contains
1272
    // [size, offset_ptr, key_data_ptr, val_data_ptr, k_nullmap_ptr, v_nullmap_pr]
1273
    // which converted results from olap_map_convertor and later will use a structure to replace it
1274
27.2k
    auto data_ptr = reinterpret_cast<const uint64_t*>(*ptr);
1275
    // total number length
1276
27.2k
    size_t element_cnt = size_t((unsigned long)(*data_ptr));
1277
27.2k
    auto offset_data = *(data_ptr + 1);
1278
27.2k
    const uint8_t* offsets_ptr = (const uint8_t*)offset_data;
1279
1280
27.2k
    if (element_cnt > 0) {
1281
43.9k
        for (size_t i = 0; i < 2; ++i) {
1282
29.3k
            auto data = *(data_ptr + 2 + i);
1283
29.3k
            auto nested_null_map = *(data_ptr + 2 + 2 + i);
1284
29.3k
            RETURN_IF_ERROR(
1285
29.3k
                    _kv_writers[i]->append(reinterpret_cast<const uint8_t*>(nested_null_map),
1286
29.3k
                                           reinterpret_cast<const void*>(data), element_cnt));
1287
29.3k
        }
1288
14.6k
    }
1289
    // make sure the order : offset writer flush next_array_item_ordinal after kv_writers append_data
1290
    // because we use _kv_writers[0]->get_next_rowid() to set next_array_item_ordinal in offset page footer
1291
27.2k
    RETURN_IF_ERROR(_offsets_writer->append_data(&offsets_ptr, num_rows));
1292
27.2k
    return Status::OK();
1293
27.2k
}
1294
1295
26.9k
Status MapColumnWriter::write_data() {
1296
26.9k
    RETURN_IF_ERROR(_offsets_writer->write_data());
1297
26.9k
    if (is_nullable()) {
1298
9.81k
        RETURN_IF_ERROR(_null_writer->write_data());
1299
9.81k
    }
1300
53.8k
    for (auto& sub_writer : _kv_writers) {
1301
53.8k
        RETURN_IF_ERROR(sub_writer->write_data());
1302
53.8k
    }
1303
26.9k
    return Status::OK();
1304
26.9k
}
1305
1306
26.6k
Status MapColumnWriter::write_ordinal_index() {
1307
26.6k
    RETURN_IF_ERROR(_offsets_writer->write_ordinal_index());
1308
26.6k
    if (is_nullable()) {
1309
9.55k
        RETURN_IF_ERROR(_null_writer->write_ordinal_index());
1310
9.55k
    }
1311
53.2k
    for (auto& sub_writer : _kv_writers) {
1312
53.2k
        if (sub_writer->get_next_rowid() != 0) {
1313
28.1k
            RETURN_IF_ERROR(sub_writer->write_ordinal_index());
1314
28.1k
        }
1315
53.2k
    }
1316
26.6k
    return Status::OK();
1317
26.6k
}
1318
1319
0
Status MapColumnWriter::append_nulls(size_t num_rows) {
1320
0
    for (auto& sub_writer : _kv_writers) {
1321
0
        RETURN_IF_ERROR(sub_writer->append_nulls(num_rows));
1322
0
    }
1323
0
    const UInt64 offset = cast_set<UInt64>(_kv_writers[0]->get_next_rowid());
1324
0
    std::vector<UInt64> offsets_data(num_rows + 1, offset);
1325
0
    const uint8_t* offsets_ptr = reinterpret_cast<const uint8_t*>(offsets_data.data());
1326
0
    RETURN_IF_ERROR(_offsets_writer->append_data(&offsets_ptr, num_rows));
1327
1328
0
    if (is_nullable()) {
1329
0
        std::vector<UInt8> null_signs(num_rows, 1);
1330
0
        const uint8_t* null_sign_ptr = null_signs.data();
1331
0
        RETURN_IF_ERROR(_null_writer->append_data(&null_sign_ptr, num_rows));
1332
0
    }
1333
0
    return Status::OK();
1334
0
}
1335
1336
0
Status MapColumnWriter::finish_current_page() {
1337
0
    return Status::NotSupported("map writer has no data, can not finish_current_page");
1338
0
}
1339
1340
11.4k
Status MapColumnWriter::write_inverted_index() {
1341
11.4k
    if (_opts.need_inverted_index) {
1342
0
        return _index_builder->finish();
1343
0
    }
1344
11.4k
    return Status::OK();
1345
11.4k
}
1346
1347
VariantColumnWriter::VariantColumnWriter(const ColumnWriterOptions& opts, TabletColumnPtr column)
1348
6.04k
        : ColumnWriter(std::move(column), opts.meta->is_nullable(), opts.meta) {
1349
6.04k
    _impl = std::make_unique<VariantColumnWriterImpl>(opts, get_column());
1350
6.04k
}
1351
1352
6.06k
Status VariantColumnWriter::init() {
1353
6.06k
    return _impl->init();
1354
6.06k
}
1355
1356
815
Status VariantColumnWriter::append_data(const uint8_t** ptr, size_t num_rows) {
1357
815
    _next_rowid += num_rows;
1358
815
    return _impl->append_data(ptr, num_rows);
1359
815
}
1360
1361
1.04k
uint64_t VariantColumnWriter::estimate_buffer_size() {
1362
1.04k
    return _impl->estimate_buffer_size();
1363
1.04k
}
1364
1365
6.09k
Status VariantColumnWriter::finish() {
1366
6.09k
    return _impl->finish();
1367
6.09k
}
1368
6.09k
Status VariantColumnWriter::write_data() {
1369
6.09k
    return _impl->write_data();
1370
6.09k
}
1371
6.09k
Status VariantColumnWriter::write_ordinal_index() {
1372
6.09k
    return _impl->write_ordinal_index();
1373
6.09k
}
1374
1375
6.08k
Status VariantColumnWriter::write_zone_map() {
1376
6.08k
    return _impl->write_zone_map();
1377
6.08k
}
1378
1379
6.06k
Status VariantColumnWriter::write_inverted_index() {
1380
6.06k
    return _impl->write_inverted_index();
1381
6.06k
}
1382
6.06k
Status VariantColumnWriter::write_bloom_filter_index() {
1383
6.06k
    return _impl->write_bloom_filter_index();
1384
6.06k
}
1385
1386
Status VariantColumnWriter::append_nullable(const uint8_t* null_map, const uint8_t** ptr,
1387
5.50k
                                            size_t num_rows) {
1388
5.50k
    return _impl->append_nullable(null_map, ptr, num_rows);
1389
5.50k
}
1390
1391
} // namespace doris::segment_v2