Coverage Report

Created: 2026-10-09 15:03

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/load/routine_load/data_consumer.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 <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