Coverage Report

Created: 2026-06-01 12:33

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