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 <aws/kinesis/KinesisClient.h> |
21 | | #include <aws/kinesis/model/GetRecordsRequest.h> |
22 | | #include <aws/kinesis/model/GetRecordsResult.h> |
23 | | #include <aws/kinesis/model/GetShardIteratorRequest.h> |
24 | | #include <aws/kinesis/model/ListShardsRequest.h> |
25 | | #include <aws/kinesis/model/Record.h> |
26 | | #include <stdint.h> |
27 | | |
28 | | #include <ctime> |
29 | | #include <functional> |
30 | | #include <map> |
31 | | #include <memory> |
32 | | #include <mutex> |
33 | | #include <ostream> |
34 | | #include <set> |
35 | | #include <string> |
36 | | #include <unordered_map> |
37 | | #include <vector> |
38 | | |
39 | | #include "common/logging.h" |
40 | | #include "common/status.h" |
41 | | #include "librdkafka/rdkafkacpp.h" |
42 | | #include "load/routine_load/kinesis_conf.h" |
43 | | #include "load/stream_load/stream_load_context.h" |
44 | | #include "runtime/aws_msk_iam_auth.h" |
45 | | #include "util/uid_util.h" |
46 | | |
47 | | namespace doris { |
48 | | |
49 | | template <typename T> |
50 | | class BlockingQueue; |
51 | | |
52 | | struct KinesisQueueItem { |
53 | | std::string shard_id; |
54 | | std::shared_ptr<Aws::Kinesis::Model::Record> record; |
55 | | bool end_of_shard = false; |
56 | | std::map<std::string, std::set<std::string>> child_shard_parent_ids; |
57 | | }; |
58 | | |
59 | | class DataConsumer { |
60 | | public: |
61 | | DataConsumer() |
62 | 26 | : _id(UniqueId::gen_uid()), |
63 | 26 | _grp_id(UniqueId::gen_uid()), |
64 | 26 | _has_grp(false), |
65 | 26 | _init(false), |
66 | 26 | _cancelled(false), |
67 | 26 | _last_visit_time(0) {} |
68 | | |
69 | 26 | virtual ~DataConsumer() {} |
70 | | |
71 | | // init the consumer with the given parameters |
72 | | virtual Status init(std::shared_ptr<StreamLoadContext> ctx) = 0; |
73 | | // start consuming |
74 | | virtual Status consume(std::shared_ptr<StreamLoadContext> ctx) = 0; |
75 | | // cancel the consuming process. |
76 | | // if the consumer is not initialized, or the consuming |
77 | | // process is already finished, call cancel() will |
78 | | // return ERROR |
79 | | virtual Status cancel(std::shared_ptr<StreamLoadContext> ctx) = 0; |
80 | | // reset the data consumer before being reused |
81 | | virtual Status reset() = 0; |
82 | | // return true the if the consumer match the need |
83 | | virtual bool match(std::shared_ptr<StreamLoadContext> ctx) = 0; |
84 | | |
85 | 0 | const UniqueId& id() { return _id; } |
86 | 0 | time_t last_visit_time() { return _last_visit_time; } |
87 | 10 | void set_grp(const UniqueId& grp_id) { |
88 | 10 | _grp_id = grp_id; |
89 | 10 | _has_grp = true; |
90 | 10 | } |
91 | | |
92 | | protected: |
93 | | UniqueId _id; |
94 | | UniqueId _grp_id; |
95 | | bool _has_grp; |
96 | | |
97 | | // lock to protect the following bools |
98 | | std::mutex _lock; |
99 | | bool _init; |
100 | | bool _cancelled; |
101 | | time_t _last_visit_time; |
102 | | }; |
103 | | |
104 | | class PShardInfo; |
105 | | |
106 | | class PIntegerPair; |
107 | | |
108 | | class KafkaEventCb : public RdKafka::EventCb { |
109 | | public: |
110 | 19 | void event_cb(RdKafka::Event& event) { |
111 | 19 | switch (event.type()) { |
112 | 7 | case RdKafka::Event::EVENT_ERROR: |
113 | 7 | LOG(INFO) << "kafka error: " << RdKafka::err2str(event.err()) |
114 | 7 | << ", event: " << event.str(); |
115 | 7 | break; |
116 | 0 | case RdKafka::Event::EVENT_STATS: |
117 | 0 | LOG(INFO) << "kafka stats: " << event.str(); |
118 | 0 | break; |
119 | | |
120 | 12 | case RdKafka::Event::EVENT_LOG: |
121 | 12 | LOG(INFO) << "kafka log-" << event.severity() << "-" << event.fac().c_str() |
122 | 12 | << ", event: " << event.str(); |
123 | 12 | break; |
124 | | |
125 | 0 | case RdKafka::Event::EVENT_THROTTLE: |
126 | 0 | LOG(INFO) << "kafka throttled: " << event.throttle_time() << "ms by " |
127 | 0 | << event.broker_name() << " id " << event.broker_id(); |
128 | 0 | break; |
129 | | |
130 | 0 | default: |
131 | 0 | LOG(INFO) << "kafka event: " << event.type() |
132 | 0 | << ", err: " << RdKafka::err2str(event.err()) << ", event: " << event.str(); |
133 | 0 | break; |
134 | 19 | } |
135 | 19 | } |
136 | | }; |
137 | | |
138 | | class KafkaDataConsumer : public DataConsumer { |
139 | | public: |
140 | | KafkaDataConsumer(std::shared_ptr<StreamLoadContext> ctx) |
141 | 2 | : _brokers(ctx->kafka_info->brokers), _topic(ctx->kafka_info->topic) {} |
142 | | |
143 | 2 | virtual ~KafkaDataConsumer() { |
144 | 2 | VLOG_NOTICE << "deconstruct consumer"; |
145 | 2 | if (_k_consumer) { |
146 | 2 | _k_consumer->close(); |
147 | 2 | delete _k_consumer; |
148 | 2 | _k_consumer = nullptr; |
149 | 2 | } |
150 | 2 | } |
151 | | |
152 | | Status init(std::shared_ptr<StreamLoadContext> ctx) override; |
153 | | // TODO(cmy): currently do not implement single consumer start method, using group_consume |
154 | 0 | Status consume(std::shared_ptr<StreamLoadContext> ctx) override { return Status::OK(); } |
155 | | Status cancel(std::shared_ptr<StreamLoadContext> ctx) override; |
156 | | // reassign partition topics |
157 | | virtual Status reset() override; |
158 | | bool match(std::shared_ptr<StreamLoadContext> ctx) override; |
159 | | // commit kafka offset |
160 | | Status commit(std::vector<RdKafka::TopicPartition*>& offset); |
161 | | |
162 | | Status assign_topic_partitions(const std::map<int32_t, int64_t>& begin_partition_offset, |
163 | | const std::string& topic, |
164 | | std::shared_ptr<StreamLoadContext> ctx); |
165 | | |
166 | | // start the consumer and put msgs to queue |
167 | | Status group_consume(BlockingQueue<RdKafka::Message*>* queue, int64_t max_running_time_ms); |
168 | | |
169 | | // get the partitions ids of the topic |
170 | | Status get_partition_meta(std::vector<int32_t>* partition_ids); |
171 | | // get offsets for times |
172 | | Status get_offsets_for_times(const std::vector<PIntegerPair>& times, |
173 | | std::vector<PIntegerPair>* offsets, int timeout); |
174 | | // get latest offsets for partitions |
175 | | Status get_latest_offsets_for_partitions(const std::vector<int32_t>& partition_ids, |
176 | | std::vector<PIntegerPair>* offsets, int timeout); |
177 | | // get offsets for times |
178 | | Status get_real_offsets_for_partitions(const std::vector<PIntegerPair>& offset_flags, |
179 | | std::vector<PIntegerPair>* offsets, int timeout); |
180 | | |
181 | | private: |
182 | | std::string _brokers; |
183 | | std::string _topic; |
184 | | std::unordered_map<std::string, std::string> _custom_properties; |
185 | | std::set<int32_t> _consuming_partition_ids; |
186 | | |
187 | | KafkaEventCb _k_event_cb; |
188 | | RdKafka::KafkaConsumer* _k_consumer = nullptr; |
189 | | |
190 | | // AWS MSK IAM authentication callback (must outlive _k_consumer) |
191 | | std::unique_ptr<AwsMskIamOAuthCallback> _aws_msk_oauth_callback; |
192 | | }; |
193 | | |
194 | | // AWS Kinesis Data Consumer |
195 | | // Consumes data from AWS Kinesis Data Streams for routine load jobs. |
196 | | // Kinesis is similar to Kafka but uses shards instead of partitions |
197 | | // and sequence numbers (strings) instead of offsets (integers). |
198 | | class KinesisDataConsumer : public DataConsumer { |
199 | | public: |
200 | | KinesisDataConsumer(std::shared_ptr<StreamLoadContext> ctx, int scan_request_timeout_ms = 0); |
201 | | virtual ~KinesisDataConsumer(); |
202 | | |
203 | | // DataConsumer interface implementation |
204 | | Status init(std::shared_ptr<StreamLoadContext> ctx) override; |
205 | 0 | Status consume(std::shared_ptr<StreamLoadContext> ctx) override { return Status::OK(); } |
206 | | Status cancel(std::shared_ptr<StreamLoadContext> ctx) override; |
207 | | Status reset() override; |
208 | | bool match(std::shared_ptr<StreamLoadContext> ctx) override; |
209 | | |
210 | | // Kinesis-specific methods |
211 | | // Assign shards with their starting sequence numbers |
212 | | Status assign_shards(const std::map<std::string, std::string>& shard_sequence_numbers, |
213 | | const std::string& stream_name, std::shared_ptr<StreamLoadContext> ctx); |
214 | | |
215 | | // Main consumption loop - pulls records from all assigned shards |
216 | | Status group_consume(BlockingQueue<KinesisQueueItem>* queue, int64_t max_running_time_ms); |
217 | | |
218 | | // Get list of shard IDs |
219 | | Status get_shard_list(std::vector<PShardInfo>* shard_infos); |
220 | | |
221 | | // Resolve LATEST by scanning retained records without loading them into Doris. |
222 | | Status get_latest_sequence_number(const std::string& shard_id, |
223 | | const std::function<Status()>& check_status, |
224 | | std::string* sequence_number); |
225 | | |
226 | | private: |
227 | | // Configuration - Basic AWS settings |
228 | | // Nonzero only for dedicated metadata scan consumers. |
229 | | const int _scan_request_timeout_ms; |
230 | | std::string _region; |
231 | | std::string _stream; |
232 | | std::string _endpoint; // Optional custom endpoint (e.g., LocalStack) |
233 | | |
234 | | // Type 1: Doris-internal parameters (not passed to AWS SDK) |
235 | | std::unordered_map<std::string, std::string> _doris_internal_properties; |
236 | | |
237 | | // Type 2: Frequently-used AWS parameters (explicit members for performance) |
238 | | // These are parsed from aws.kinesis.* properties during init() |
239 | | std::vector<std::string> _explicit_shards; // aws.kinesis.shards (comma-separated) |
240 | | std::string _default_position; // aws.kinesis.default.pos (LATEST/TRIM_HORIZON) |
241 | | std::map<std::string, std::string> |
242 | | _shard_positions; // aws.kinesis.shards.pos (shard_id:position) |
243 | | |
244 | | // Type 3: Less-frequently-used AWS API parameters (wrapped in KinesisConf) |
245 | | std::unique_ptr<KinesisConf> _kinesis_conf; |
246 | | |
247 | | // AWS credentials and other properties |
248 | | std::unordered_map<std::string, std::string> _custom_properties; |
249 | | |
250 | | // Active shards being consumed |
251 | | std::set<std::string> _consuming_shard_ids; |
252 | | |
253 | | // AWS Kinesis client |
254 | | std::shared_ptr<Aws::Kinesis::KinesisClient> _kinesis_client; |
255 | | |
256 | | // Shard iterator management |
257 | | // Kinesis requires shard iterators to consume records |
258 | | // shard_id -> current shard iterator |
259 | | std::map<std::string, std::string> _shard_iterators; |
260 | | |
261 | | // Tracks the MillisBehindLatest value per shard from the last GetRecords call. |
262 | | // Updated during group_consume; read by the task executor to populate ctx after consumption. |
263 | | std::map<std::string, int64_t> _millis_behind_latest; |
264 | | |
265 | | public: |
266 | | // Returns the MillisBehindLatest snapshot collected during group_consume. |
267 | 8 | const std::map<std::string, int64_t>& get_millis_behind_latest() const { |
268 | 8 | return _millis_behind_latest; |
269 | 8 | } |
270 | | |
271 | | private: |
272 | | // Helper methods |
273 | | // Create and configure AWS Kinesis client with credentials |
274 | | Status _create_kinesis_client(std::shared_ptr<StreamLoadContext> ctx); |
275 | | |
276 | | // Get shard iterator for a shard at a specific sequence number position |
277 | | Status _get_shard_iterator(const std::string& shard_id, const std::string& sequence_number, |
278 | | std::string* iterator); |
279 | | |
280 | | enum class EnqueueResult { COMPLETE, QUEUE_SHUTDOWN }; |
281 | | |
282 | | // Queue shutdown is normal batch completion, but the caller must stop the whole consumer. |
283 | | EnqueueResult _process_records(const std::string& shard_id, |
284 | | Aws::Kinesis::Model::GetRecordsResult result, |
285 | | BlockingQueue<KinesisQueueItem>* queue, int64_t* received_rows, |
286 | | int64_t* put_rows); |
287 | | |
288 | | // Check if an AWS error is retriable (throttling, network, etc.) |
289 | | bool _is_retriable_error(const Aws::Client::AWSError<Aws::Kinesis::KinesisErrors>& error); |
290 | | }; |
291 | | |
292 | | } // end namespace doris |