Coverage Report

Created: 2026-10-02 07:29

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/storage/segment/column_writer.h
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
#pragma once
19
20
#include <gen_cpp/AgentService_types.h>
21
#include <gen_cpp/olap_file.pb.h>
22
#include <gen_cpp/segment_v2.pb.h>
23
#include <stddef.h>
24
#include <stdint.h>
25
26
#include <algorithm>
27
#include <memory> // for unique_ptr
28
#include <ostream>
29
#include <span>
30
#include <string>
31
#include <unordered_map>
32
#include <utility>
33
#include <vector>
34
35
#include "common/status.h" // for Status
36
#include "storage/index/ann/ann_index_writer.h"
37
#include "storage/index/bloom_filter/bloom_filter.h"
38
#include "storage/index/inverted/inverted_index_writer.h"
39
#include "storage/segment/common.h"
40
#include "storage/segment/options.h"
41
#include "storage/segment/variant/nested_group_provider.h"
42
#include "storage/segment/variant/variant_statistics.h"
43
#include "storage/tablet/tablet_schema.h" // for TabletColumnPtr
44
#include "storage/types.h"                // for field_type_size
45
#include "util/bitmap.h"                  // for BitmapChange
46
#include "util/slice.h"                   // for OwnedSlice
47
48
namespace doris {
49
50
class BlockCompressionCodec;
51
class TabletColumn;
52
class TabletIndex;
53
struct RowsetWriterContext;
54
struct VariantColumnData;
55
56
namespace io {
57
class FileWriter;
58
}
59
60
namespace segment_v2 {
61
62
enum class VariantWriterInputFormat : uint8_t {
63
    UNSET,
64
    V2,
65
};
66
67
struct ColumnWriterOptions {
68
    // input and output parameter:
69
    // - input: column_id/unique_id/type/length/encoding/compression/is_nullable members
70
    // - output: encoding/indexes/dict_page members
71
    ColumnMetaPB* meta = nullptr;
72
    size_t data_page_size = STORAGE_PAGE_SIZE_DEFAULT_VALUE;
73
    size_t dict_page_size = STORAGE_DICT_PAGE_SIZE_DEFAULT_VALUE;
74
    // store compressed page only when space saving is above the threshold.
75
    // space saving = 1 - compressed_size / uncompressed_size
76
    double compression_min_space_saving = 0.1;
77
    bool need_zone_map = false;
78
    bool need_bloom_filter = false;
79
    bool is_ngram_bf_index = false;
80
    bool need_inverted_index = false;
81
    bool need_ann_index = false;
82
    uint8_t gram_size;
83
    uint16_t gram_bf_size;
84
    BloomFilterOptions bf_options;
85
    std::vector<const TabletIndex*> inverted_indexes;
86
    IndexFileWriter* index_file_writer = nullptr;
87
    // The owning segment serves a direct load (stream/broker load,
88
    // DataWriteType::TYPE_DIRECT) rather than compaction / schema change. Set
89
    // once by the segment writer and propagated to variant subcolumn writers;
90
    // forwarded to every created IndexColumnWriter via set_direct_load() so
91
    // SNII can select its direct-load PRX zstd level without plumbing
92
    // DataWriteType itself down here.
93
    bool is_direct_load = false;
94
95
    SegmentFooterPB* footer = nullptr;
96
    io::FileWriter* file_writer = nullptr;
97
    CompressionTypePB compression_type = UNKNOWN_COMPRESSION;
98
    RowsetWriterContext* rowset_ctx = nullptr;
99
    // For collect segment statistics for compaction
100
    std::vector<RowsetReaderSharedPtr> input_rs_readers;
101
    const TabletIndex* ann_index = nullptr;
102
103
    // Storage format of the owning tablet (V2 or V3). Set once by the segment writer
104
    // (from TabletMeta::storage_format()) and propagated down to aux child writers
105
    // (null / array-length / map-length), struct subcolumn writers and variant subcolumn
106
    // writers. All encoding-default decisions consult this via resolve_default_encoding().
107
    // Also forwarded to BinaryDictPageBuilder via PageBuilderOptions::binary_plain_encoding.
108
    TabletStorageFormatPB storage_format = TabletStorageFormatPB::TABLET_STORAGE_FORMAT_V2;
109
110
0
    std::string to_string() const {
111
0
        std::stringstream ss;
112
0
        ss << std::boolalpha << "meta=" << meta->DebugString()
113
0
           << ", data_page_size=" << data_page_size << ", dict_page_size=" << dict_page_size
114
0
           << ", compression_min_space_saving = " << compression_min_space_saving
115
0
           << ", need_zone_map=" << need_zone_map << ", need_bloom_filter" << need_bloom_filter;
116
0
        return ss.str();
117
0
    }
118
};
119
120
class EncodingInfo;
121
class NullBitmapBuilder;
122
class OrdinalIndexWriter;
123
class PageBuilder;
124
class BloomFilterIndexWriter;
125
class ZoneMapIndexWriter;
126
class VariantColumnWriterImpl;
127
class VariantShredder;
128
class VariantPathBuilder;
129
class ColumnWriter;
130
131
class ColumnWriter {
132
public:
133
    static Status create(const ColumnWriterOptions& opts, const TabletColumn* column,
134
                         io::FileWriter* file_writer, std::unique_ptr<ColumnWriter>* writer);
135
    static Status create_struct_writer(const ColumnWriterOptions& opts, const TabletColumn* column,
136
                                       io::FileWriter* file_writer,
137
                                       std::unique_ptr<ColumnWriter>* writer);
138
    static Status create_array_writer(const ColumnWriterOptions& opts, const TabletColumn* column,
139
                                      io::FileWriter* file_writer,
140
                                      std::unique_ptr<ColumnWriter>* writer);
141
    static Status create_map_writer(const ColumnWriterOptions& opts, const TabletColumn* column,
142
                                    io::FileWriter* file_writer,
143
                                    std::unique_ptr<ColumnWriter>* writer);
144
145
    static Status create_variant_writer(const ColumnWriterOptions& opts, const TabletColumn* column,
146
                                        io::FileWriter* file_writer,
147
                                        std::unique_ptr<ColumnWriter>* writer);
148
149
    static Status create_agg_state_writer(const ColumnWriterOptions& opts,
150
                                          const TabletColumn* column, io::FileWriter* file_writer,
151
                                          std::unique_ptr<ColumnWriter>* writer);
152
153
    explicit ColumnWriter(TabletColumnPtr column, bool is_nullable, ColumnMetaPB* meta);
154
155
1.06M
    virtual ~ColumnWriter() = default;
156
157
    virtual Status init() = 0;
158
159
    template <typename CellType>
160
    Status append(const CellType& cell) {
161
        if (_is_nullable) {
162
            uint8_t nullmap = 0;
163
            BitmapChange(&nullmap, 0, cell.is_null());
164
            return append_nullable(&nullmap, cell.cell_ptr(), 1);
165
        } else {
166
            auto* cel_ptr = cell.cell_ptr();
167
            return append_data((const uint8_t**)&cel_ptr, 1);
168
        }
169
    }
170
171
    // Now we only support append one by one, we should support append
172
    // multi rows in one call
173
    Status append(bool is_null, void* data) {
174
        uint8_t nullmap = 0;
175
        BitmapChange(&nullmap, 0, is_null);
176
        return append_nullable(&nullmap, data, 1);
177
    }
178
179
    Status append(const uint8_t* nullmap, const void* data, size_t num_rows);
180
181
    Status append_nullable(const uint8_t* nullmap, const void* data, size_t num_rows);
182
183
    // use only in vectorized load
184
    virtual Status append_nullable(const uint8_t* null_map, const uint8_t** data, size_t num_rows);
185
186
    virtual Status append_nulls(size_t num_rows) = 0;
187
188
    virtual Status finish_current_page() = 0;
189
190
    virtual uint64_t estimate_buffer_size() = 0;
191
192
    // finish append data
193
    virtual Status finish() = 0;
194
195
    // write all data into file
196
    virtual Status write_data() = 0;
197
198
    virtual Status write_ordinal_index() = 0;
199
200
    virtual Status write_zone_map() = 0;
201
202
    virtual Status write_inverted_index() = 0;
203
204
685k
    virtual Status write_ann_index() { return Status::OK(); }
205
206
    virtual Status write_bloom_filter_index() = 0;
207
208
    virtual ordinal_t get_next_rowid() const = 0;
209
210
    virtual uint64_t get_raw_data_bytes() const = 0;
211
    virtual uint64_t get_total_uncompressed_data_pages_bytes() const = 0;
212
    virtual uint64_t get_total_compressed_data_pages_bytes() const = 0;
213
214
    // used for append not null data.
215
    virtual Status append_data(const uint8_t** ptr, size_t num_rows) = 0;
216
217
5.26M
    bool is_nullable() const { return _is_nullable; }
218
219
2.57M
    const TabletColumn* get_column() const { return _column.get(); }
220
221
    // Per-row in-memory cell footprint of this writer's column, used to step
222
    // the input pointer across rows in append_*/null-run loops.
223
3.41M
    size_t cell_size() const { return field_type_size(_column->type()); }
224
225
725k
    ColumnMetaPB* get_column_meta() const { return _column_meta; }
226
227
protected:
228
    DataTypePtr _data_type;
229
230
private:
231
    TabletColumnPtr _column;
232
    bool _is_nullable;
233
    ColumnMetaPB* _column_meta;
234
};
235
236
class FlushPageCallback {
237
public:
238
66.6k
    virtual ~FlushPageCallback() = default;
239
0
    virtual void put_extra_info_in_page(DataPageFooterPB* footer) {}
240
};
241
242
// Encode one column's data into some memory slice.
243
// Because some columns would be stored in a file, we should wait
244
// until all columns has been finished, and then data can be written
245
// to file
246
class ScalarColumnWriter : public ColumnWriter {
247
public:
248
    ScalarColumnWriter(const ColumnWriterOptions& opts, TabletColumnPtr column,
249
                       io::FileWriter* file_writer);
250
251
    ~ScalarColumnWriter() override;
252
253
    Status init() override;
254
255
    Status append_nulls(size_t num_rows) override;
256
257
    Status finish_current_page() override;
258
259
    uint64_t estimate_buffer_size() override;
260
261
    // finish append data
262
    Status finish() override;
263
264
    Status write_data() override;
265
    Status write_ordinal_index() override;
266
    Status write_zone_map() override;
267
    Status write_inverted_index() override;
268
    Status write_bloom_filter_index() override;
269
159k
    ordinal_t get_next_rowid() const override { return _next_rowid; }
270
271
837k
    uint64_t get_raw_data_bytes() const override { return _raw_data_bytes; }
272
273
837k
    uint64_t get_total_uncompressed_data_pages_bytes() const override {
274
837k
        return _total_uncompressed_data_pages_size;
275
837k
    }
276
277
837k
    uint64_t get_total_compressed_data_pages_bytes() const override {
278
837k
        return _total_compressed_data_pages_size;
279
837k
    }
280
281
66.7k
    void register_flush_page_callback(FlushPageCallback* flush_page_callback) {
282
66.7k
        _new_page_callback = flush_page_callback;
283
66.7k
    }
284
    Status append_data(const uint8_t** ptr, size_t num_rows) override;
285
    Status append_nullable(const uint8_t* null_map, const uint8_t** ptr, size_t num_rows) override;
286
287
    // used for append not null data. When page is full, will append data not reach num_rows.
288
    Status append_data_in_current_page(const uint8_t** ptr, size_t* num_written);
289
290
3.13M
    Status append_data_in_current_page(const uint8_t* ptr, size_t* num_written) {
291
3.13M
        RETURN_IF_CATCH_EXCEPTION(
292
3.13M
                { return _internal_append_data_in_current_page(ptr, num_written); });
293
3.13M
    }
294
    friend class ArrayColumnWriter;
295
    friend class OffsetColumnWriter;
296
297
private:
298
    Status _internal_append_data_in_current_page(const uint8_t* ptr, size_t* num_written);
299
300
private:
301
    struct NullRun {
302
        bool is_null;
303
        uint32_t len;
304
    };
305
306
    std::vector<NullRun> _null_run_buffer;
307
    std::unique_ptr<PageBuilder> _page_builder;
308
309
    std::unique_ptr<NullBitmapBuilder> _null_bitmap_builder;
310
311
    ColumnWriterOptions _opts;
312
313
    const EncodingInfo* _encoding_info = nullptr;
314
315
    ordinal_t _next_rowid = 0;
316
317
    // All Pages will be organized into a linked list
318
    struct Page {
319
        // the data vector may contain:
320
        //     1. one OwnedSlice if the page body is compressed
321
        //     2. one OwnedSlice if the page body is not compressed and doesn't have nullmap
322
        //     3. two OwnedSlice if the page body is not compressed and has nullmap
323
        // use vector for easier management for lifetime of OwnedSlice
324
        std::vector<OwnedSlice> data;
325
        PageFooterPB footer;
326
    };
327
328
1.02M
    void _push_back_page(std::unique_ptr<Page> page) {
329
1.94M
        for (auto& data_slice : page->data) {
330
1.94M
            _data_size += data_slice.slice().size;
331
1.94M
        }
332
        // estimate (page footer + footer size + checksum) took 20 bytes
333
1.02M
        _data_size += 20;
334
        // add page to pages' tail
335
1.02M
        _pages.emplace_back(std::move(page));
336
1.02M
    }
337
338
    Status _write_data_page(Page* page);
339
340
private:
341
    io::FileWriter* _file_writer = nullptr;
342
    // total size of data page list
343
    uint64_t _data_size;
344
345
    uint64_t _raw_data_bytes {0};
346
    uint64_t _total_uncompressed_data_pages_size {0};
347
    uint64_t _total_compressed_data_pages_size {0};
348
349
    // cached generated pages,
350
    std::vector<std::unique_ptr<Page>> _pages;
351
    ordinal_t _first_rowid = 0;
352
353
    BlockCompressionCodec* _compress_codec;
354
355
    std::unique_ptr<OrdinalIndexWriter> _ordinal_index_builder;
356
    std::unique_ptr<ZoneMapIndexWriter> _zone_map_index_builder;
357
    std::vector<std::unique_ptr<IndexColumnWriter>> _inverted_index_builders;
358
    std::unique_ptr<BloomFilterIndexWriter> _bloom_filter_index_builder;
359
360
    // call before flush data page.
361
    FlushPageCallback* _new_page_callback = nullptr;
362
};
363
364
// offsetColumnWriter is used column which has offset column, like array, map.
365
//  column type is only uint64 and should response for whole column value [start, end], end will set
366
//  in footer.next_array_item_ordinal which in finish_cur_page() callback put_extra_info_in_page()
367
class OffsetColumnWriter final : public ScalarColumnWriter, FlushPageCallback {
368
public:
369
    OffsetColumnWriter(const ColumnWriterOptions& opts, TabletColumnPtr column,
370
                       io::FileWriter* file_writer);
371
372
    ~OffsetColumnWriter() override;
373
374
    Status init() override;
375
376
    Status append_data(const uint8_t** ptr, size_t num_rows) override;
377
378
private:
379
    void put_extra_info_in_page(DataPageFooterPB* footer) override;
380
381
    uint64_t _next_offset;
382
};
383
384
class StructColumnWriter final : public ColumnWriter {
385
public:
386
    explicit StructColumnWriter(const ColumnWriterOptions& opts, TabletColumnPtr column,
387
                                ScalarColumnWriter* null_writer,
388
                                std::vector<std::unique_ptr<ColumnWriter>>& sub_column_writers);
389
2.98k
    ~StructColumnWriter() override = default;
390
391
    Status init() override;
392
393
    Status append_nullable(const uint8_t* null_map, const uint8_t** data, size_t num_rows) override;
394
    Status append_data(const uint8_t** ptr, size_t num_rows) override;
395
396
    uint64_t estimate_buffer_size() override;
397
398
    Status finish() override;
399
    Status write_data() override;
400
    Status write_ordinal_index() override;
401
0
    Status append_nulls(size_t num_rows) override {
402
0
        return Status::NotSupported("struct writer can not append_nulls");
403
0
    }
404
405
    Status finish_current_page() override;
406
407
2.27k
    Status write_zone_map() override {
408
2.27k
        if (_opts.need_zone_map) {
409
0
            return Status::NotSupported("struct not support zone map");
410
0
        }
411
2.27k
        return Status::OK();
412
2.27k
    }
413
414
    Status write_inverted_index() override;
415
2.27k
    Status write_bloom_filter_index() override {
416
2.27k
        if (_opts.need_bloom_filter) {
417
0
            return Status::NotSupported("struct not support bloom filter index");
418
0
        }
419
2.27k
        return Status::OK();
420
2.27k
    }
421
422
3.39k
    ordinal_t get_next_rowid() const override { return _sub_column_writers[0]->get_next_rowid(); }
423
424
2.98k
    uint64_t get_raw_data_bytes() const override {
425
2.98k
        return _get_total_data_pages_bytes(&ColumnWriter::get_raw_data_bytes);
426
2.98k
    }
427
428
2.98k
    uint64_t get_total_uncompressed_data_pages_bytes() const override {
429
2.98k
        return _get_total_data_pages_bytes(&ColumnWriter::get_total_uncompressed_data_pages_bytes);
430
2.98k
    }
431
432
2.98k
    uint64_t get_total_compressed_data_pages_bytes() const override {
433
2.98k
        return _get_total_data_pages_bytes(&ColumnWriter::get_total_compressed_data_pages_bytes);
434
2.98k
    }
435
436
private:
437
    template <typename Func>
438
8.96k
    uint64_t _get_total_data_pages_bytes(Func func) const {
439
8.96k
        uint64_t size = is_nullable() ? std::invoke(func, _null_writer.get()) : 0;
440
55.7k
        for (const auto& writer : _sub_column_writers) {
441
55.7k
            size += std::invoke(func, writer.get());
442
55.7k
        }
443
8.96k
        return size;
444
8.96k
    }
445
446
private:
447
    size_t _num_sub_column_writers;
448
    std::unique_ptr<ScalarColumnWriter> _null_writer;
449
    std::vector<std::unique_ptr<ColumnWriter>> _sub_column_writers;
450
    ColumnWriterOptions _opts;
451
};
452
453
class ArrayColumnWriter final : public ColumnWriter {
454
public:
455
    explicit ArrayColumnWriter(const ColumnWriterOptions& opts, TabletColumnPtr column,
456
                               OffsetColumnWriter* offset_writer, ScalarColumnWriter* null_writer,
457
                               std::unique_ptr<ColumnWriter> item_writer);
458
42.9k
    ~ArrayColumnWriter() override = default;
459
460
    Status init() override;
461
462
    Status append_data(const uint8_t** ptr, size_t num_rows) override;
463
464
    uint64_t estimate_buffer_size() override;
465
466
    Status finish() override;
467
    Status write_data() override;
468
    Status write_ordinal_index() override;
469
    Status append_nulls(size_t num_rows) override;
470
    Status append_nullable(const uint8_t* null_map, const uint8_t** ptr, size_t num_rows) override;
471
472
    Status finish_current_page() override;
473
474
40.1k
    Status write_zone_map() override {
475
40.1k
        if (_opts.need_zone_map) {
476
0
            return Status::NotSupported("array not support zone map");
477
0
        }
478
40.1k
        return Status::OK();
479
40.1k
    }
480
481
    Status write_inverted_index() override;
482
    Status write_ann_index() override;
483
40.1k
    Status write_bloom_filter_index() override {
484
40.1k
        if (_opts.need_bloom_filter) {
485
0
            return Status::NotSupported("array not support bloom filter index");
486
0
        }
487
40.1k
        return Status::OK();
488
40.1k
    }
489
44.8k
    ordinal_t get_next_rowid() const override { return _offset_writer->get_next_rowid(); }
490
491
42.2k
    uint64_t get_raw_data_bytes() const override {
492
42.2k
        return _get_total_data_pages_bytes(&ColumnWriter::get_raw_data_bytes);
493
42.2k
    }
494
495
42.2k
    uint64_t get_total_uncompressed_data_pages_bytes() const override {
496
42.2k
        return _get_total_data_pages_bytes(&ColumnWriter::get_total_uncompressed_data_pages_bytes);
497
42.2k
    }
498
499
42.3k
    uint64_t get_total_compressed_data_pages_bytes() const override {
500
42.3k
        return _get_total_data_pages_bytes(&ColumnWriter::get_total_compressed_data_pages_bytes);
501
42.3k
    }
502
503
private:
504
    template <typename Func>
505
126k
    uint64_t _get_total_data_pages_bytes(Func func) const {
506
126k
        uint64_t size = std::invoke(func, _offset_writer.get());
507
126k
        if (is_nullable()) {
508
78.4k
            size += std::invoke(func, _null_writer.get());
509
78.4k
        }
510
126k
        size += std::invoke(func, _item_writer.get());
511
126k
        return size;
512
126k
    }
513
514
private:
515
    Status write_null_column(size_t num_rows, bool is_null); // 写入num_rows个null标记
516
42.5k
    bool has_empty_items() const { return _item_writer->get_next_rowid() == 0; }
517
518
private:
519
    std::unique_ptr<OffsetColumnWriter> _offset_writer;
520
    std::unique_ptr<ScalarColumnWriter> _null_writer;
521
    std::unique_ptr<ColumnWriter> _item_writer;
522
    std::unique_ptr<IndexColumnWriter> _inverted_index_writer;
523
    std::unique_ptr<AnnIndexColumnWriter> _ann_index_writer;
524
    ColumnWriterOptions _opts;
525
};
526
527
class MapColumnWriter final : public ColumnWriter {
528
public:
529
    explicit MapColumnWriter(const ColumnWriterOptions& opts, TabletColumnPtr column,
530
                             ScalarColumnWriter* null_writer, OffsetColumnWriter* offsets_writer,
531
                             std::vector<std::unique_ptr<ColumnWriter>>& _kv_writers);
532
533
23.6k
    ~MapColumnWriter() override = default;
534
535
    Status init() override;
536
537
    Status append_data(const uint8_t** ptr, size_t num_rows) override;
538
    Status append_nullable(const uint8_t* null_map, const uint8_t** ptr, size_t num_rows) override;
539
    uint64_t estimate_buffer_size() override;
540
541
    Status finish() override;
542
    Status write_data() override;
543
    Status write_ordinal_index() override;
544
    Status write_inverted_index() override;
545
0
    Status append_nulls(size_t num_rows) override {
546
0
        return Status::NotSupported("map writer can not append_nulls");
547
0
    }
548
549
    Status finish_current_page() override;
550
551
10.0k
    Status write_zone_map() override {
552
10.0k
        if (_opts.need_zone_map) {
553
0
            return Status::NotSupported("map not support zone map");
554
0
        }
555
10.0k
        return Status::OK();
556
10.0k
    }
557
558
10.0k
    Status write_bloom_filter_index() override {
559
10.0k
        if (_opts.need_bloom_filter) {
560
0
            return Status::NotSupported("map not support bloom filter index");
561
0
        }
562
10.0k
        return Status::OK();
563
10.0k
    }
564
565
    // according key writer to get next rowid
566
24.6k
    ordinal_t get_next_rowid() const override { return _offsets_writer->get_next_rowid(); }
567
568
11.0k
    uint64_t get_raw_data_bytes() const override {
569
11.0k
        return _get_total_data_pages_bytes(&ColumnWriter::get_raw_data_bytes);
570
11.0k
    }
571
572
11.0k
    uint64_t get_total_uncompressed_data_pages_bytes() const override {
573
11.0k
        return _get_total_data_pages_bytes(&ColumnWriter::get_total_uncompressed_data_pages_bytes);
574
11.0k
    }
575
576
11.0k
    uint64_t get_total_compressed_data_pages_bytes() const override {
577
11.0k
        return _get_total_data_pages_bytes(&ColumnWriter::get_total_compressed_data_pages_bytes);
578
11.0k
    }
579
580
private:
581
    template <typename Func>
582
33.1k
    uint64_t _get_total_data_pages_bytes(Func func) const {
583
33.1k
        uint64_t size = std::invoke(func, _offsets_writer.get());
584
33.1k
        if (is_nullable()) {
585
25.8k
            size += std::invoke(func, _null_writer.get());
586
25.8k
        }
587
66.2k
        for (const auto& writer : _kv_writers) {
588
66.2k
            size += std::invoke(func, writer.get());
589
66.2k
        }
590
33.1k
        return size;
591
33.1k
    }
592
593
private:
594
    std::vector<std::unique_ptr<ColumnWriter>> _kv_writers;
595
    // we need null writer to make sure a row is null or not
596
    std::unique_ptr<ScalarColumnWriter> _null_writer;
597
    std::unique_ptr<OffsetColumnWriter> _offsets_writer;
598
    std::unique_ptr<IndexColumnWriter> _index_builder;
599
    ColumnWriterOptions _opts;
600
};
601
602
// used for compaction to write sub variant column
603
class VariantSubcolumnWriter : public ColumnWriter {
604
public:
605
    explicit VariantSubcolumnWriter(const ColumnWriterOptions& opts, TabletColumnPtr column);
606
607
    ~VariantSubcolumnWriter() override;
608
609
    Status init() override;
610
611
    Status append_data(const uint8_t** ptr, size_t num_rows) override;
612
613
    uint64_t estimate_buffer_size() override;
614
615
    Status finish() override;
616
    Status write_data() override;
617
    Status write_ordinal_index() override;
618
619
    Status write_zone_map() override;
620
621
    Status write_inverted_index() override;
622
    Status write_bloom_filter_index() override;
623
2
    ordinal_t get_next_rowid() const override { return _next_rowid; }
624
625
22
    uint64_t get_raw_data_bytes() const override {
626
22
        return 0; // TODO
627
22
    }
628
629
22
    uint64_t get_total_uncompressed_data_pages_bytes() const override {
630
22
        return 0; // TODO
631
22
    }
632
633
22
    uint64_t get_total_compressed_data_pages_bytes() const override {
634
22
        return 0; // TODO
635
22
    }
636
637
0
    Status append_nulls(size_t num_rows) override {
638
0
        return Status::NotSupported("variant writer can not append_nulls");
639
0
    }
640
    Status append_nullable(const uint8_t* null_map, const uint8_t** ptr, size_t num_rows) override;
641
642
0
    Status finish_current_page() override {
643
0
        return Status::NotSupported("variant writer has no data, can not finish_current_page");
644
0
    }
645
646
0
    size_t get_non_null_size() const { return none_null_size; }
647
648
    Status finalize();
649
650
private:
651
    Status _append(const uint8_t* null_map, const uint8_t** ptr, size_t num_rows);
652
    Status _append_v2(const VariantColumnData& column, size_t num_rows,
653
                      std::span<const uint8_t> outer_nulls);
654
    Status _ensure_input_format(const VariantColumnData& column);
655
    Status _initialize_v2_builder();
656
    bool is_finalized() const;
657
    bool _is_finalized = false;
658
    ordinal_t _next_rowid = 0;
659
    size_t none_null_size = 0;
660
    VariantWriterInputFormat _input_format = VariantWriterInputFormat::UNSET;
661
    std::unique_ptr<VariantPathBuilder> _v2_builder;
662
    size_t _num_rows = 0;
663
    ColumnWriterOptions _opts;
664
    std::unique_ptr<ColumnWriter> _writer;
665
    TabletIndexes _indexes;
666
667
    std::unique_ptr<NestedGroupWriteProvider> _nested_group_provider;
668
    VariantStatistics _statistics;
669
};
670
671
class VariantColumnWriter : public ColumnWriter {
672
public:
673
    explicit VariantColumnWriter(const ColumnWriterOptions& opts, TabletColumnPtr column);
674
675
5.36k
    ~VariantColumnWriter() override = default;
676
677
    Status init() override;
678
679
    Status append_data(const uint8_t** ptr, size_t num_rows) override;
680
681
    uint64_t estimate_buffer_size() override;
682
683
    Status finish() override;
684
    Status write_data() override;
685
    Status write_ordinal_index() override;
686
687
    Status write_zone_map() override;
688
689
    Status write_inverted_index() override;
690
    Status write_bloom_filter_index() override;
691
1
    ordinal_t get_next_rowid() const override { return _next_rowid; }
692
693
5.32k
    uint64_t get_raw_data_bytes() const override {
694
5.32k
        return 0; // TODO
695
5.32k
    }
696
697
5.32k
    uint64_t get_total_uncompressed_data_pages_bytes() const override {
698
5.32k
        return 0; // TODO
699
5.32k
    }
700
701
5.32k
    uint64_t get_total_compressed_data_pages_bytes() const override {
702
5.32k
        return 0; // TODO
703
5.32k
    }
704
705
0
    Status append_nulls(size_t num_rows) override {
706
0
        return Status::NotSupported("variant writer can not append_nulls");
707
0
    }
708
    Status append_nullable(const uint8_t* null_map, const uint8_t** ptr, size_t num_rows) override;
709
710
0
    Status finish_current_page() override {
711
0
        return Status::NotSupported("variant writer has no data, can not finish_current_page");
712
0
    }
713
714
    VariantColumnWriterImpl* impl_for_test() const { return _impl.get(); }
715
716
private:
717
    std::unique_ptr<VariantColumnWriterImpl> _impl;
718
    ordinal_t _next_rowid = 0;
719
};
720
721
} // namespace segment_v2
722
} // namespace doris