/src/duckdb/extension/parquet/writer/struct_column_writer.cpp
Line | Count | Source |
1 | | #include <stdint.h> |
2 | | #include <string> |
3 | | #include <utility> |
4 | | #include <vector> |
5 | | |
6 | | #include "duckdb/common/vector/struct_vector.hpp" |
7 | | #include "writer/struct_column_writer.hpp" |
8 | | #include "column_writer.hpp" |
9 | | #include "duckdb/common/helper.hpp" |
10 | | #include "duckdb/common/numeric_utils.hpp" |
11 | | #include "duckdb/common/optional_idx.hpp" |
12 | | #include "duckdb/common/typedefs.hpp" |
13 | | #include "duckdb/common/types/vector.hpp" |
14 | | #include "duckdb/common/unique_ptr.hpp" |
15 | | #include "duckdb/common/vector.hpp" |
16 | | #include "duckdb/common/vector/flat_vector.hpp" |
17 | | #include "parquet_column_schema.hpp" |
18 | | #include "parquet_types.h" |
19 | | |
20 | | namespace duckdb { |
21 | | |
22 | | using namespace duckdb_parquet; // NOLINT |
23 | | |
24 | | using duckdb_parquet::ConvertedType; |
25 | | using duckdb_parquet::FieldRepetitionType; |
26 | | |
27 | | class StructColumnWriterState : public ColumnWriterState { |
28 | | public: |
29 | | StructColumnWriterState(duckdb_parquet::RowGroup &row_group, idx_t col_idx) |
30 | 0 | : row_group(row_group), col_idx(col_idx) { |
31 | 0 | } |
32 | 0 | ~StructColumnWriterState() override = default; |
33 | | |
34 | | duckdb_parquet::RowGroup &row_group; |
35 | | idx_t col_idx; |
36 | | vector<unique_ptr<ColumnWriterState>> child_states; |
37 | | }; |
38 | | |
39 | 0 | unique_ptr<ColumnWriterState> StructColumnWriter::InitializeWriteState(duckdb_parquet::RowGroup &row_group) { |
40 | 0 | auto result = make_uniq<StructColumnWriterState>(row_group, row_group.columns.size()); |
41 | |
|
42 | 0 | result->child_states.reserve(child_writers.size()); |
43 | 0 | for (auto &child_writer : child_writers) { |
44 | 0 | result->child_states.push_back(child_writer->InitializeWriteState(row_group)); |
45 | 0 | } |
46 | 0 | return std::move(result); |
47 | 0 | } |
48 | | |
49 | 0 | bool StructColumnWriter::HasAnalyze() { |
50 | 0 | for (auto &child_writer : child_writers) { |
51 | 0 | if (child_writer->HasAnalyze()) { |
52 | 0 | return true; |
53 | 0 | } |
54 | 0 | } |
55 | 0 | return false; |
56 | 0 | } |
57 | | |
58 | 0 | void StructColumnWriter::Analyze(ColumnWriterState &state_p, ColumnWriterState *parent, Vector &vector, idx_t count) { |
59 | 0 | auto &state = state_p.Cast<StructColumnWriterState>(); |
60 | 0 | auto &child_vectors = StructVector::GetEntries(vector); |
61 | 0 | for (idx_t child_idx = 0; child_idx < child_writers.size(); child_idx++) { |
62 | | // Need to check again. It might be that just one child needs it but the rest not |
63 | 0 | if (child_writers[child_idx]->HasAnalyze()) { |
64 | 0 | child_writers[child_idx]->Analyze(*state.child_states[child_idx], &state_p, child_vectors[child_idx], |
65 | 0 | count); |
66 | 0 | } |
67 | 0 | } |
68 | 0 | } |
69 | | |
70 | 0 | void StructColumnWriter::FinalizeAnalyze(ColumnWriterState &state_p) { |
71 | 0 | auto &state = state_p.Cast<StructColumnWriterState>(); |
72 | 0 | for (idx_t child_idx = 0; child_idx < child_writers.size(); child_idx++) { |
73 | | // Need to check again. It might be that just one child needs it but the rest not |
74 | 0 | if (child_writers[child_idx]->HasAnalyze()) { |
75 | 0 | child_writers[child_idx]->FinalizeAnalyze(*state.child_states[child_idx]); |
76 | 0 | } |
77 | 0 | } |
78 | 0 | } |
79 | | |
80 | | void StructColumnWriter::Prepare(ColumnWriterState &state_p, ColumnWriterState *parent, Vector &vector, idx_t count, |
81 | 0 | bool vector_can_span_multiple_pages) { |
82 | 0 | auto &state = state_p.Cast<StructColumnWriterState>(); |
83 | |
|
84 | 0 | auto &validity = FlatVector::ValidityMutable(vector); |
85 | 0 | if (parent) { |
86 | | // propagate empty entries from the parent |
87 | 0 | if (state.is_empty.size() < parent->is_empty.size()) { |
88 | 0 | state.is_empty.insert(state.is_empty.end(), |
89 | 0 | parent->is_empty.begin() + |
90 | 0 | NumericCast<duckdb::vector<bool>::difference_type>(state.is_empty.size()), |
91 | 0 | parent->is_empty.end()); |
92 | 0 | } |
93 | 0 | } |
94 | 0 | HandleRepeatLevels(state_p, parent, count); |
95 | 0 | HandleDefineLevels(state_p, parent, validity, count, PARQUET_DEFINE_VALID, MaxDefine() - 1); |
96 | 0 | auto &child_vectors = StructVector::GetEntries(vector); |
97 | 0 | for (idx_t child_idx = 0; child_idx < child_writers.size(); child_idx++) { |
98 | 0 | child_writers[child_idx]->Prepare(*state.child_states[child_idx], &state_p, child_vectors[child_idx], count, |
99 | 0 | vector_can_span_multiple_pages); |
100 | 0 | } |
101 | 0 | } |
102 | | |
103 | 0 | void StructColumnWriter::BeginWrite(ColumnWriterState &state_p) { |
104 | 0 | auto &state = state_p.Cast<StructColumnWriterState>(); |
105 | 0 | for (idx_t child_idx = 0; child_idx < child_writers.size(); child_idx++) { |
106 | 0 | child_writers[child_idx]->BeginWrite(*state.child_states[child_idx]); |
107 | 0 | } |
108 | 0 | } |
109 | | |
110 | 0 | void StructColumnWriter::Write(ColumnWriterState &state_p, Vector &vector, idx_t count) { |
111 | 0 | auto &state = state_p.Cast<StructColumnWriterState>(); |
112 | 0 | auto &child_vectors = StructVector::GetEntries(vector); |
113 | 0 | for (idx_t child_idx = 0; child_idx < child_writers.size(); child_idx++) { |
114 | 0 | child_writers[child_idx]->Write(*state.child_states[child_idx], child_vectors[child_idx], count); |
115 | 0 | } |
116 | 0 | } |
117 | | |
118 | 0 | void StructColumnWriter::PrepareWrite(ColumnWriterState &state_p) { |
119 | 0 | auto &state = state_p.Cast<StructColumnWriterState>(); |
120 | 0 | for (idx_t child_idx = 0; child_idx < child_writers.size(); child_idx++) { |
121 | | // we add the null count of the struct to the null count of the children |
122 | 0 | state.child_states[child_idx]->null_count += state_p.null_count; |
123 | 0 | child_writers[child_idx]->PrepareWrite(*state.child_states[child_idx]); |
124 | 0 | } |
125 | 0 | } |
126 | | |
127 | 0 | void StructColumnWriter::FinalizeWrite(ColumnWriterState &state_p) { |
128 | 0 | auto &state = state_p.Cast<StructColumnWriterState>(); |
129 | 0 | for (idx_t child_idx = 0; child_idx < child_writers.size(); child_idx++) { |
130 | 0 | child_writers[child_idx]->FinalizeWrite(*state.child_states[child_idx]); |
131 | 0 | } |
132 | 0 | } |
133 | | |
134 | 0 | idx_t StructColumnWriter::FinalizeSchema(vector<duckdb_parquet::SchemaElement> &schemas) { |
135 | 0 | idx_t schema_idx = schemas.size(); |
136 | |
|
137 | 0 | auto &schema = column_schema; |
138 | 0 | schema.SetSchemaIndex(schema_idx); |
139 | |
|
140 | 0 | auto &repetition_type = schema.repetition_type; |
141 | 0 | auto &name = schema.name; |
142 | 0 | auto &field_id = schema.field_id; |
143 | | |
144 | | // set up the schema element for this struct |
145 | 0 | duckdb_parquet::SchemaElement schema_element; |
146 | 0 | schema_element.repetition_type = repetition_type; |
147 | 0 | schema_element.num_children = NumericCast<int32_t>(child_writers.size()); |
148 | 0 | schema_element.__isset.num_children = true; |
149 | 0 | schema_element.__isset.type = false; |
150 | 0 | schema_element.__isset.repetition_type = true; |
151 | 0 | schema_element.name = name; |
152 | 0 | if (field_id.IsValid()) { |
153 | 0 | schema_element.__isset.field_id = true; |
154 | 0 | schema_element.field_id = NumericCast<int32_t>(field_id.GetIndex()); |
155 | 0 | } |
156 | 0 | schemas.push_back(std::move(schema_element)); |
157 | |
|
158 | 0 | idx_t unique_columns = 0; |
159 | 0 | for (auto &child_writer : child_writers) { |
160 | 0 | unique_columns += child_writer->FinalizeSchema(schemas); |
161 | 0 | } |
162 | 0 | return unique_columns; |
163 | 0 | } |
164 | | |
165 | | } // namespace duckdb |