Coverage Report

Created: 2026-08-28 15:07

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/load/channel/load_stream.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 "load/channel/load_stream.h"
19
20
#include <brpc/stream.h>
21
#include <bthread/bthread.h>
22
#include <bthread/condition_variable.h>
23
#include <bthread/mutex.h>
24
25
#include <memory>
26
#include <sstream>
27
28
#include "bvar/bvar.h"
29
#include "cloud/config.h"
30
#include "common/signal_handler.h"
31
#include "load/channel/load_channel.h"
32
#include "load/channel/load_stream_mgr.h"
33
#include "load/channel/load_stream_writer.h"
34
#include "load/delta_writer/delta_writer.h"
35
#include "runtime/exec_env.h"
36
#include "runtime/fragment_mgr.h"
37
#include "runtime/runtime_profile.h"
38
#include "runtime/workload_group/workload_group_manager.h"
39
#include "storage/rowset/rowset_factory.h"
40
#include "storage/rowset/rowset_meta.h"
41
#include "storage/storage_engine.h"
42
#include "storage/tablet/tablet.h"
43
#include "storage/tablet/tablet_fwd.h"
44
#include "storage/tablet/tablet_manager.h"
45
#include "storage/tablet/tablet_schema.h"
46
#include "storage/tablet_info.h"
47
#include "util/debug_points.h"
48
#include "util/thrift_util.h"
49
#include "util/uid_util.h"
50
51
#define UNKNOWN_ID_FOR_TEST 0x7c00
52
53
namespace doris {
54
55
bvar::Adder<int64_t> g_load_stream_cnt("load_stream_count");
56
bvar::LatencyRecorder g_load_stream_flush_wait_ms("load_stream_flush_wait_ms");
57
bvar::Adder<int> g_load_stream_flush_running_threads("load_stream_flush_wait_threads");
58
59
TabletStream::TabletStream(const PUniqueId& load_id, int64_t id, int64_t txn_id,
60
                           LoadStreamMgr* load_stream_mgr, RuntimeProfile* profile)
61
561
        : _id(id),
62
561
          _next_segid(0),
63
561
          _load_id(load_id),
64
561
          _txn_id(txn_id),
65
561
          _load_stream_mgr(load_stream_mgr) {
66
561
    load_stream_mgr->create_token(_flush_token);
67
561
    _profile = profile->create_child(fmt::format("TabletStream {}", id), true, true);
68
561
    _append_data_timer = ADD_TIMER(_profile, "AppendDataTime");
69
561
    _add_segment_timer = ADD_TIMER(_profile, "AddSegmentTime");
70
561
    _close_wait_timer = ADD_TIMER(_profile, "CloseWaitTime");
71
561
}
72
73
10
inline std::ostream& operator<<(std::ostream& ostr, const TabletStream& tablet_stream) {
74
10
    ostr << "load_id=" << print_id(tablet_stream._load_id) << ", txn_id=" << tablet_stream._txn_id
75
10
         << ", tablet_id=" << tablet_stream._id << ", status=" << tablet_stream._status.status();
76
10
    return ostr;
77
10
}
78
79
Status TabletStream::init(std::shared_ptr<OlapTableSchemaParam> schema, int64_t index_id,
80
561
                          int64_t partition_id) {
81
561
    WriteRequest req {
82
561
            .tablet_id = _id,
83
561
            .txn_id = _txn_id,
84
561
            .index_id = index_id,
85
561
            .partition_id = partition_id,
86
561
            .load_id = _load_id,
87
561
            .table_schema_param = schema,
88
            // TODO(plat1ko): write_file_cache
89
561
            .storage_vault_id {},
90
561
    };
91
92
561
    _load_stream_writer = std::make_shared<LoadStreamWriter>(&req, _profile);
93
561
    DBUG_EXECUTE_IF("TabletStream.init.uninited_writer", {
94
561
        _status.update(Status::Uninitialized("fault injection"));
95
561
        return _status.status();
96
561
    });
97
561
    _status.update(_load_stream_writer->init());
98
561
    if (!_status.ok()) {
99
1
        LOG(INFO) << "failed to init rowset builder due to " << *this;
100
1
    }
101
561
    return _status.status();
102
561
}
103
104
11.1k
Status TabletStream::append_data(const PStreamHeader& header, butil::IOBuf* data) {
105
11.1k
    if (!_status.ok()) {
106
1
        return _status.status();
107
1
    }
108
109
    // dispatch add_segment request
110
11.1k
    if (header.opcode() == PStreamHeader::ADD_SEGMENT) {
111
168
        return add_segment(header, data);
112
168
    }
113
114
10.9k
    SCOPED_TIMER(_append_data_timer);
115
116
10.9k
    int64_t src_id = header.src_id();
117
10.9k
    uint32_t segid = header.segment_id();
118
    // Ensure there are enough space and mapping are built.
119
10.9k
    SegIdMapping* mapping = nullptr;
120
10.9k
    {
121
10.9k
        std::lock_guard lock_guard(_lock);
122
10.9k
        if (!_segids_mapping.contains(src_id)) {
123
182
            _segids_mapping[src_id] = std::make_unique<SegIdMapping>();
124
182
        }
125
10.9k
        mapping = _segids_mapping[src_id].get();
126
10.9k
    }
127
10.9k
    if (segid + 1 > mapping->size()) {
128
        // TODO: Each sender lock is enough.
129
182
        std::lock_guard lock_guard(_lock);
130
182
        ssize_t origin_size = mapping->size();
131
183
        if (segid + 1 > origin_size) {
132
183
            mapping->resize(segid + 1, std::numeric_limits<uint32_t>::max());
133
375
            for (size_t index = origin_size; index <= segid; index++) {
134
192
                mapping->at(index) = _next_segid;
135
192
                _next_segid++;
136
192
                VLOG_DEBUG << "src_id=" << src_id << ", segid=" << index << " to "
137
0
                           << " segid=" << _next_segid - 1 << ", " << *this;
138
192
            }
139
183
        }
140
182
    }
141
142
    // Each sender sends data in one segment sequential, so we also do not
143
    // need a lock here.
144
10.9k
    bool eos = header.segment_eos();
145
10.9k
    FileType file_type = header.file_type();
146
10.9k
    uint32_t new_segid = mapping->at(segid);
147
10.9k
    DCHECK(new_segid != std::numeric_limits<uint32_t>::max());
148
10.9k
    butil::IOBuf buf = data->movable();
149
10.9k
    auto flush_func = [this, new_segid, eos, buf, header, file_type]() mutable {
150
10.9k
        signal::set_signal_task_id(_load_id);
151
10.9k
        g_load_stream_flush_running_threads << -1;
152
10.9k
        auto st = _load_stream_writer->append_data(new_segid, header.offset(), buf, file_type);
153
10.9k
        if (!st.ok() && !config::is_cloud_mode()) {
154
1
            auto res = ExecEnv::get_tablet(_id);
155
1
            TabletSharedPtr tablet =
156
1
                    res.has_value() ? std::dynamic_pointer_cast<Tablet>(res.value()) : nullptr;
157
1
            if (tablet) {
158
1
                tablet->report_error(st);
159
1
            }
160
1
        }
161
10.9k
        if (eos && st.ok()) {
162
189
            DBUG_EXECUTE_IF("TabletStream.append_data.unknown_file_type",
163
189
                            { file_type = static_cast<FileType>(-1); });
164
189
            if (file_type == FileType::SEGMENT_FILE || file_type == FileType::INVERTED_INDEX_FILE) {
165
189
                st = _load_stream_writer->close_writer(new_segid, file_type);
166
189
            } else {
167
0
                st = Status::InternalError(
168
0
                        "appent data failed, file type error, file type = {}, "
169
0
                        "segment_id={}",
170
0
                        file_type, new_segid);
171
0
            }
172
189
        }
173
10.9k
        DBUG_EXECUTE_IF("TabletStream.append_data.append_failed",
174
10.9k
                        { st = Status::InternalError("fault injection"); });
175
10.9k
        if (!st.ok()) {
176
2
            _status.update(st);
177
2
            LOG(WARNING) << "write data failed " << st << ", " << *this;
178
2
        }
179
10.9k
    };
180
10.9k
    auto load_stream_flush_token_max_tasks = config::load_stream_flush_token_max_tasks;
181
10.9k
    auto load_stream_max_wait_flush_token_time_ms =
182
10.9k
            config::load_stream_max_wait_flush_token_time_ms;
183
10.9k
    DBUG_EXECUTE_IF("TabletStream.append_data.long_wait", {
184
10.9k
        load_stream_flush_token_max_tasks = 0;
185
10.9k
        load_stream_max_wait_flush_token_time_ms = 1000;
186
10.9k
    });
187
10.9k
    MonotonicStopWatch timer;
188
10.9k
    timer.start();
189
10.9k
    while (_flush_token->num_tasks() >= load_stream_flush_token_max_tasks) {
190
20
        if (timer.elapsed_time() / 1000 / 1000 >= load_stream_max_wait_flush_token_time_ms) {
191
0
            _status.update(
192
0
                    Status::Error<true>("wait flush token back pressure time is more than "
193
0
                                        "load_stream_max_wait_flush_token_time {}",
194
0
                                        load_stream_max_wait_flush_token_time_ms));
195
0
            return _status.status();
196
0
        }
197
20
        bthread_usleep(2 * 1000); // 2ms
198
20
    }
199
10.9k
    timer.stop();
200
10.9k
    int64_t time_ms = timer.elapsed_time() / 1000 / 1000;
201
10.9k
    g_load_stream_flush_wait_ms << time_ms;
202
10.9k
    g_load_stream_flush_running_threads << 1;
203
10.9k
    Status st = Status::OK();
204
10.9k
    DBUG_EXECUTE_IF("TabletStream.append_data.submit_func_failed",
205
10.9k
                    { st = Status::InternalError("fault injection"); });
206
10.9k
    if (st.ok()) {
207
10.9k
        st = _flush_token->submit_func(flush_func);
208
10.9k
    }
209
10.9k
    if (!st.ok()) {
210
0
        _status.update(st);
211
0
    }
212
10.9k
    return _status.status();
213
10.9k
}
214
215
168
Status TabletStream::add_segment(const PStreamHeader& header, butil::IOBuf* data) {
216
168
    if (!_status.ok()) {
217
0
        return _status.status();
218
0
    }
219
220
168
    SCOPED_TIMER(_add_segment_timer);
221
168
    DCHECK(header.has_segment_statistics());
222
168
    SegmentStatistics stat(header.segment_statistics());
223
224
168
    int64_t src_id = header.src_id();
225
168
    uint32_t segid = header.segment_id();
226
168
    uint32_t new_segid;
227
168
    DBUG_EXECUTE_IF("TabletStream.add_segment.unknown_segid", { segid = UNKNOWN_ID_FOR_TEST; });
228
168
    {
229
168
        std::lock_guard lock_guard(_lock);
230
168
        if (!_segids_mapping.contains(src_id)) {
231
0
            _status.update(Status::InternalError(
232
0
                    "add segment failed, no segment written by this src be yet, src_id={}, "
233
0
                    "segment_id={}",
234
0
                    src_id, segid));
235
0
            return _status.status();
236
0
        }
237
168
        DBUG_EXECUTE_IF("TabletStream.add_segment.segid_never_written",
238
168
                        { segid = static_cast<uint32_t>(_segids_mapping[src_id]->size()); });
239
168
        if (segid >= _segids_mapping[src_id]->size()) {
240
0
            _status.update(Status::InternalError(
241
0
                    "add segment failed, segment is never written, src_id={}, segment_id={}",
242
0
                    src_id, segid));
243
0
            return _status.status();
244
0
        }
245
168
        new_segid = _segids_mapping[src_id]->at(segid);
246
168
    }
247
168
    DCHECK(new_segid != std::numeric_limits<uint32_t>::max());
248
249
168
    auto add_segment_func = [this, new_segid, stat]() {
250
168
        signal::set_signal_task_id(_load_id);
251
168
        auto st = _load_stream_writer->add_segment(new_segid, stat);
252
168
        DBUG_EXECUTE_IF("TabletStream.add_segment.add_segment_failed",
253
168
                        { st = Status::InternalError("fault injection"); });
254
168
        if (!st.ok()) {
255
0
            _status.update(st);
256
0
            LOG(INFO) << "add segment failed " << *this;
257
0
        }
258
168
    };
259
168
    Status st = Status::OK();
260
168
    DBUG_EXECUTE_IF("TabletStream.add_segment.submit_func_failed",
261
168
                    { st = Status::InternalError("fault injection"); });
262
168
    if (st.ok()) {
263
168
        st = _flush_token->submit_func(add_segment_func);
264
168
    }
265
168
    if (!st.ok()) {
266
0
        _status.update(st);
267
0
    }
268
168
    return _status.status();
269
168
}
270
271
1.67k
Status TabletStream::_run_in_heavy_work_pool(std::function<Status()> fn) {
272
1.67k
    bthread::Mutex mu;
273
1.67k
    std::unique_lock<bthread::Mutex> lock(mu);
274
1.67k
    bthread::ConditionVariable cv;
275
1.67k
    auto st = Status::OK();
276
1.67k
    auto func = [this, &mu, &cv, &st, &fn] {
277
1.67k
        signal::set_signal_task_id(_load_id);
278
1.67k
        st = fn();
279
1.67k
        std::lock_guard<bthread::Mutex> lock(mu);
280
1.67k
        cv.notify_one();
281
1.67k
    };
282
1.67k
    bool ret = _load_stream_mgr->heavy_work_pool()->try_offer(func);
283
1.67k
    if (!ret) {
284
0
        return Status::Error<ErrorCode::INTERNAL_ERROR>(
285
0
                "there is not enough thread resource for close load");
286
0
    }
287
1.67k
    cv.wait(lock);
288
1.67k
    return st;
289
1.67k
}
290
291
1.12k
void TabletStream::wait_for_flush_tasks() {
292
1.12k
    {
293
1.12k
        std::lock_guard lock_guard(_lock);
294
1.12k
        if (_flush_tasks_done) {
295
561
            return;
296
561
        }
297
561
        _flush_tasks_done = true;
298
561
    }
299
300
561
    if (!_status.ok()) {
301
2
        _flush_token->shutdown();
302
2
        return;
303
2
    }
304
305
    // Use heavy_work_pool to avoid blocking bthread
306
559
    auto st = _run_in_heavy_work_pool([this]() {
307
559
        _flush_token->wait();
308
559
        return Status::OK();
309
559
    });
310
559
    if (!st.ok()) {
311
        // If heavy_work_pool is unavailable, fall back to shutdown
312
        // which will cancel pending tasks and wait for running tasks
313
0
        _flush_token->shutdown();
314
0
        _status.update(st);
315
0
    }
316
559
}
317
318
561
void TabletStream::pre_close() {
319
561
    SCOPED_TIMER(_close_wait_timer);
320
561
    wait_for_flush_tasks();
321
322
561
    if (!_status.ok()) {
323
3
        return;
324
3
    }
325
326
558
    DBUG_EXECUTE_IF("TabletStream.close.segment_num_mismatch", { _num_segments++; });
327
558
    if (_check_num_segments && (_next_segid.load() != _num_segments)) {
328
2
        _status.update(Status::Corruption(
329
2
                "segment num mismatch in tablet {}, expected: {}, actual: {}, load_id: {}", _id,
330
2
                _num_segments, _next_segid.load(), print_id(_load_id)));
331
2
        return;
332
2
    }
333
334
556
    _status.update(_run_in_heavy_work_pool([this]() { return _load_stream_writer->pre_close(); }));
335
556
}
336
337
561
Status TabletStream::close() {
338
561
    if (!_status.ok()) {
339
6
        return _status.status();
340
6
    }
341
342
555
    SCOPED_TIMER(_close_wait_timer);
343
555
    _status.update(_run_in_heavy_work_pool([this]() { return _load_stream_writer->close(); }));
344
555
    return _status.status();
345
561
}
346
347
IndexStream::IndexStream(const PUniqueId& load_id, int64_t id, int64_t txn_id,
348
                         std::shared_ptr<OlapTableSchemaParam> schema,
349
                         LoadStreamMgr* load_stream_mgr, RuntimeProfile* profile)
350
108
        : _id(id),
351
108
          _load_id(load_id),
352
108
          _txn_id(txn_id),
353
108
          _schema(schema),
354
108
          _load_stream_mgr(load_stream_mgr) {
355
108
    _profile = profile->create_child(fmt::format("IndexStream {}", id), true, true);
356
108
    _append_data_timer = ADD_TIMER(_profile, "AppendDataTime");
357
108
    _close_wait_timer = ADD_TIMER(_profile, "CloseWaitTime");
358
108
}
359
360
108
IndexStream::~IndexStream() {
361
    // Ensure all TabletStreams have their flush tokens properly handled before destruction.
362
    // In normal flow, close() should have called pre_close() on all tablet streams.
363
    // But if IndexStream is destroyed without close() being called (e.g., on_idle_timeout),
364
    // we need to wait for flush tasks here to ensure flush tokens are properly shut down.
365
561
    for (auto& [_, tablet_stream] : _tablet_streams_map) {
366
561
        tablet_stream->wait_for_flush_tasks();
367
561
    }
368
108
}
369
370
11.1k
Status IndexStream::append_data(const PStreamHeader& header, butil::IOBuf* data) {
371
11.1k
    SCOPED_TIMER(_append_data_timer);
372
11.1k
    int64_t tablet_id = header.tablet_id();
373
11.1k
    TabletStreamSharedPtr tablet_stream;
374
11.1k
    {
375
11.1k
        std::lock_guard lock_guard(_lock);
376
11.1k
        auto it = _tablet_streams_map.find(tablet_id);
377
11.1k
        if (it == _tablet_streams_map.end()) {
378
181
            _init_tablet_stream(tablet_stream, tablet_id, header.partition_id());
379
10.9k
        } else {
380
10.9k
            tablet_stream = it->second;
381
10.9k
        }
382
11.1k
    }
383
384
11.1k
    return tablet_stream->append_data(header, data);
385
11.1k
}
386
387
void IndexStream::_init_tablet_stream(TabletStreamSharedPtr& tablet_stream, int64_t tablet_id,
388
561
                                      int64_t partition_id) {
389
561
    tablet_stream = std::make_shared<TabletStream>(_load_id, tablet_id, _txn_id, _load_stream_mgr,
390
561
                                                   _profile);
391
561
    _tablet_streams_map[tablet_id] = tablet_stream;
392
561
    auto st = tablet_stream->init(_schema, _id, partition_id);
393
561
    if (!st.ok()) {
394
1
        LOG(WARNING) << "tablet stream init failed " << *tablet_stream;
395
1
    }
396
561
}
397
398
168
void IndexStream::get_all_write_tablet_ids(std::vector<int64_t>* tablet_ids) {
399
168
    std::lock_guard lock_guard(_lock);
400
506
    for (const auto& [tablet_id, _] : _tablet_streams_map) {
401
506
        tablet_ids->push_back(tablet_id);
402
506
    }
403
168
}
404
405
void IndexStream::close(const std::vector<PTabletID>& tablets_to_commit,
406
108
                        std::vector<int64_t>* success_tablet_ids, FailedTablets* failed_tablets) {
407
108
    std::lock_guard lock_guard(_lock);
408
108
    SCOPED_TIMER(_close_wait_timer);
409
    // open all need commit tablets
410
586
    for (const auto& tablet : tablets_to_commit) {
411
586
        if (_id != tablet.index_id()) {
412
21
            continue;
413
21
        }
414
565
        TabletStreamSharedPtr tablet_stream;
415
565
        auto it = _tablet_streams_map.find(tablet.tablet_id());
416
565
        if (it == _tablet_streams_map.end()) {
417
380
            _init_tablet_stream(tablet_stream, tablet.tablet_id(), tablet.partition_id());
418
380
        } else {
419
185
            tablet_stream = it->second;
420
185
        }
421
565
        if (tablet.has_num_segments()) {
422
562
            tablet_stream->add_num_segments(tablet.num_segments());
423
562
        } else {
424
            // for compatibility reasons (sink from old version BE)
425
3
            tablet_stream->disable_num_segments_check();
426
3
        }
427
565
    }
428
429
561
    for (auto& [_, tablet_stream] : _tablet_streams_map) {
430
561
        tablet_stream->pre_close();
431
561
    }
432
433
561
    for (auto& [_, tablet_stream] : _tablet_streams_map) {
434
561
        auto st = tablet_stream->close();
435
561
        if (st.ok()) {
436
555
            success_tablet_ids->push_back(tablet_stream->id());
437
555
        } else {
438
6
            LOG(INFO) << "close tablet stream " << *tablet_stream << ", status=" << st;
439
6
            failed_tablets->emplace_back(tablet_stream->id(), st);
440
6
        }
441
561
    }
442
108
}
443
444
// TODO: Profile is temporary disabled, because:
445
// 1. It's not being processed by the upstream for now
446
// 2. There are some problems in _profile->to_thrift()
447
LoadStream::LoadStream(const PUniqueId& load_id, LoadStreamMgr* load_stream_mgr,
448
                       bool enable_profile)
449
78
        : _load_id(load_id), _enable_profile(false), _load_stream_mgr(load_stream_mgr) {
450
78
    g_load_stream_cnt << 1;
451
78
    _profile = std::make_unique<RuntimeProfile>("LoadStream");
452
78
    _append_data_timer = ADD_TIMER(_profile, "AppendDataTime");
453
78
    _close_wait_timer = ADD_TIMER(_profile, "CloseWaitTime");
454
78
    TUniqueId load_tid = ((UniqueId)load_id).to_thrift();
455
78
#ifndef BE_TEST
456
78
    std::shared_ptr<QueryContext> query_context =
457
78
            ExecEnv::GetInstance()->fragment_mgr()->get_query_ctx(load_tid);
458
78
    if (query_context != nullptr) {
459
78
        _resource_ctx = query_context->resource_ctx();
460
78
    } else {
461
0
        _resource_ctx = ResourceContext::create_shared();
462
0
        _resource_ctx->task_controller()->set_task_id(load_tid);
463
0
        std::shared_ptr<MemTrackerLimiter> mem_tracker = MemTrackerLimiter::create_shared(
464
0
                MemTrackerLimiter::Type::LOAD,
465
0
                fmt::format("(FromLoadStream)Load#Id={}", ((UniqueId)load_id).to_string()));
466
0
        _resource_ctx->memory_context()->set_mem_tracker(mem_tracker);
467
0
    }
468
#else
469
    _resource_ctx = ResourceContext::create_shared();
470
    _resource_ctx->task_controller()->set_task_id(load_tid);
471
    std::shared_ptr<MemTrackerLimiter> mem_tracker = MemTrackerLimiter::create_shared(
472
            MemTrackerLimiter::Type::LOAD,
473
            fmt::format("(FromLoadStream)Load#Id={}", ((UniqueId)load_id).to_string()));
474
    _resource_ctx->memory_context()->set_mem_tracker(mem_tracker);
475
#endif
476
78
}
477
478
93
LoadStream::~LoadStream() {
479
93
    g_load_stream_cnt << -1;
480
93
    LOG(INFO) << "load stream is deconstructed " << *this;
481
93
}
482
483
93
Status LoadStream::init(const POpenLoadStreamRequest* request) {
484
93
    _txn_id = request->txn_id();
485
93
    _total_streams = static_cast<int32_t>(request->total_streams());
486
93
    _is_incremental = (_total_streams == 0);
487
488
93
    _schema = std::make_shared<OlapTableSchemaParam>();
489
93
    RETURN_IF_ERROR(_schema->init(request->schema()));
490
108
    for (auto& index : request->schema().indexes()) {
491
108
        _index_streams_map[index.id()] = std::make_shared<IndexStream>(
492
108
                _load_id, index.id(), _txn_id, _schema, _load_stream_mgr, _profile.get());
493
108
    }
494
93
    LOG(INFO) << "succeed to init load stream " << *this;
495
93
    return Status::OK();
496
93
}
497
498
bool LoadStream::close(int64_t src_id, const std::vector<PTabletID>& tablets_to_commit,
499
175
                       std::vector<int64_t>* success_tablet_ids, FailedTablets* failed_tablets) {
500
175
    std::lock_guard<bthread::Mutex> lock_guard(_lock);
501
175
    SCOPED_TIMER(_close_wait_timer);
502
503
    // we do nothing until recv CLOSE_LOAD from all stream to ensure all data are handled before ack
504
175
    _open_streams[src_id]--;
505
175
    if (_open_streams[src_id] == 0) {
506
97
        _open_streams.erase(src_id);
507
97
    }
508
175
    _close_load_cnt++;
509
175
    LOG(INFO) << "received CLOSE_LOAD from sender " << src_id << ", remaining "
510
175
              << _total_streams - _close_load_cnt << " senders, " << *this;
511
512
175
    _tablets_to_commit.insert(_tablets_to_commit.end(), tablets_to_commit.begin(),
513
175
                              tablets_to_commit.end());
514
515
175
    if (_close_load_cnt < _total_streams) {
516
        // do not return commit info if there is remaining streams.
517
82
        return false;
518
82
    }
519
520
108
    for (auto& [_, index_stream] : _index_streams_map) {
521
108
        index_stream->close(_tablets_to_commit, success_tablet_ids, failed_tablets);
522
108
    }
523
93
    LOG(INFO) << "close load " << *this << ", success_tablet_num=" << success_tablet_ids->size()
524
93
              << ", failed_tablet_num=" << failed_tablets->size();
525
93
    return true;
526
175
}
527
528
175
std::vector<int64_t> LoadStream::mark_eos_sent_and_collect(int64_t stream_id, bool is_incremental) {
529
175
    std::lock_guard<bthread::Mutex> lock_guard(_lock);
530
175
    std::vector<int64_t> to_close;
531
    // A non-incremental stream is closed as soon as its own CLOSE_LOAD (and EOS)
532
    // is handled -- this is the first batch of streams, known up front, not subject
533
    // to fencing. Closing it promptly also means a duplicate/late CLOSE_LOAD lands
534
    // on an already-closed stream and is dropped, instead of being counted again.
535
    // An incremental stream must be deferred (fencing #56120: it may only close once
536
    // every non-incremental stream is closed), so it is parked in _eos_sent_stream_ids
537
    // until all CLOSE_LOADs have been received.
538
175
    if (is_incremental) {
539
        // Parked only after the caller sent this stream's EOS via _report_result,
540
        // so every parked id is safe to close.
541
3
        _eos_sent_stream_ids.push_back(stream_id);
542
172
    } else {
543
172
        to_close.push_back(stream_id);
544
172
    }
545
    // `_close_load_cnt == _total_streams` means every CLOSE_LOAD has been counted by
546
    // close(). Latch it so that any thread reaching here afterwards also drains the
547
    // parked incremental streams, guaranteeing none is left un-closed regardless of
548
    // thread interleaving (fixes the split-lock leak race).
549
175
    if (_close_load_cnt >= _total_streams) {
550
99
        _all_close_load_received = true;
551
99
    }
552
175
    if (_all_close_load_received) {
553
99
        for (const auto& parked_id : _eos_sent_stream_ids) {
554
3
            to_close.push_back(parked_id);
555
3
        }
556
99
        _eos_sent_stream_ids.clear();
557
99
    }
558
175
    return to_close;
559
175
}
560
561
void LoadStream::_report_result(StreamId stream, const Status& status,
562
                                const std::vector<int64_t>& success_tablet_ids,
563
195
                                const FailedTablets& failed_tablets, bool eos) {
564
195
    LOG(INFO) << "report result " << *this << ", success tablet num " << success_tablet_ids.size()
565
195
              << ", failed tablet num " << failed_tablets.size();
566
195
    butil::IOBuf buf;
567
195
    PLoadStreamResponse response;
568
195
    response.set_eos(eos);
569
195
    status.to_protobuf(response.mutable_status());
570
555
    for (auto& id : success_tablet_ids) {
571
555
        response.add_success_tablet_ids(id);
572
555
    }
573
195
    for (auto& [id, st] : failed_tablets) {
574
10
        auto pb = response.add_failed_tablets();
575
10
        pb->set_id(id);
576
10
        st.to_protobuf(pb->mutable_status());
577
10
    }
578
579
195
    if (_enable_profile && _close_load_cnt == _total_streams) {
580
0
        TRuntimeProfileTree tprofile;
581
0
        ThriftSerializer ser(false, 4096);
582
0
        uint8_t* profile_buf = nullptr;
583
0
        uint32_t len = 0;
584
0
        std::unique_lock<bthread::Mutex> l(_lock);
585
586
0
        _profile->to_thrift(&tprofile);
587
0
        auto st = ser.serialize(&tprofile, &len, &profile_buf);
588
0
        if (st.ok()) {
589
0
            response.set_load_stream_profile(profile_buf, len);
590
0
        } else {
591
0
            LOG(WARNING) << "TRuntimeProfileTree serialize failed, errmsg=" << st << ", " << *this;
592
0
        }
593
0
    }
594
595
195
    buf.append(response.SerializeAsString());
596
195
    auto wst = _write_stream(stream, buf);
597
195
    if (!wst.ok()) {
598
0
        LOG(WARNING) << " report result failed with " << wst << ", " << *this;
599
0
    }
600
195
}
601
602
0
void LoadStream::_report_schema(StreamId stream, const PStreamHeader& hdr) {
603
0
    butil::IOBuf buf;
604
0
    PLoadStreamResponse response;
605
0
    Status st = Status::OK();
606
0
    for (const auto& req : hdr.tablets()) {
607
0
        BaseTabletSPtr tablet;
608
0
        if (auto res = ExecEnv::get_tablet(req.tablet_id()); res.has_value()) {
609
0
            tablet = std::move(res).value();
610
0
        } else {
611
0
            st = std::move(res).error();
612
0
            break;
613
0
        }
614
0
        auto* resp = response.add_tablet_schemas();
615
0
        resp->set_index_id(req.index_id());
616
0
        resp->set_enable_unique_key_merge_on_write(tablet->enable_unique_key_merge_on_write());
617
0
        tablet->tablet_schema()->to_schema_pb(resp->mutable_tablet_schema());
618
0
    }
619
0
    st.to_protobuf(response.mutable_status());
620
621
0
    buf.append(response.SerializeAsString());
622
0
    auto wst = _write_stream(stream, buf);
623
0
    if (!wst.ok()) {
624
0
        LOG(WARNING) << " report result failed with " << wst << ", " << *this;
625
0
    }
626
0
}
627
628
168
void LoadStream::_report_tablet_load_info(StreamId stream, int64_t index_id) {
629
168
    std::vector<int64_t> write_tablet_ids;
630
168
    auto it = _index_streams_map.find(index_id);
631
168
    if (it != _index_streams_map.end()) {
632
168
        it->second->get_all_write_tablet_ids(&write_tablet_ids);
633
168
    }
634
635
168
    if (!write_tablet_ids.empty()) {
636
168
        butil::IOBuf buf;
637
168
        PLoadStreamResponse response;
638
168
        auto* tablet_load_infos = response.mutable_tablet_load_rowset_num_infos();
639
168
        _collect_tablet_load_info_from_tablets(write_tablet_ids, tablet_load_infos);
640
168
        if (tablet_load_infos->empty()) {
641
168
            return;
642
168
        }
643
0
        buf.append(response.SerializeAsString());
644
0
        auto wst = _write_stream(stream, buf);
645
0
        if (!wst.ok()) {
646
0
            LOG(WARNING) << "report tablet load info failed with " << wst << ", " << *this;
647
0
        }
648
0
    }
649
168
}
650
651
void LoadStream::_collect_tablet_load_info_from_tablets(
652
        const std::vector<int64_t>& tablet_ids,
653
168
        google::protobuf::RepeatedPtrField<PTabletLoadRowsetInfo>* tablet_load_infos) {
654
506
    for (auto tablet_id : tablet_ids) {
655
506
        BaseTabletSPtr tablet;
656
506
        if (auto res = ExecEnv::get_tablet(tablet_id); res.has_value()) {
657
506
            tablet = std::move(res).value();
658
506
        } else {
659
0
            continue;
660
0
        }
661
506
        BaseDeltaWriter::collect_tablet_load_rowset_num_info(tablet.get(), tablet_load_infos);
662
506
    }
663
168
}
664
665
195
Status LoadStream::_write_stream(StreamId stream, butil::IOBuf& buf) {
666
195
    for (;;) {
667
195
        int ret = 0;
668
195
        DBUG_EXECUTE_IF("LoadStream._write_stream.EAGAIN", { ret = EAGAIN; });
669
195
        if (ret == 0) {
670
195
            ret = brpc::StreamWrite(stream, buf);
671
195
        }
672
195
        switch (ret) {
673
195
        case 0:
674
195
            return Status::OK();
675
0
        case EAGAIN: {
676
0
            const timespec time = butil::seconds_from_now(config::load_stream_eagain_wait_seconds);
677
0
            int wait_ret = brpc::StreamWait(stream, &time);
678
0
            if (wait_ret != 0) {
679
0
                return Status::InternalError("StreamWait failed, err={}", wait_ret);
680
0
            }
681
0
            break;
682
0
        }
683
0
        default:
684
0
            return Status::InternalError("StreamWrite failed, err={}", ret);
685
195
        }
686
195
    }
687
0
    return Status::OK();
688
195
}
689
690
11.3k
void LoadStream::_parse_header(butil::IOBuf* const message, PStreamHeader& hdr) {
691
11.3k
    butil::IOBufAsZeroCopyInputStream wrapper(*message);
692
11.3k
    hdr.ParseFromZeroCopyStream(&wrapper);
693
11.3k
    VLOG_DEBUG << "header parse result: " << hdr.DebugString();
694
11.3k
}
695
696
11.1k
Status LoadStream::_append_data(const PStreamHeader& header, butil::IOBuf* data) {
697
11.1k
    SCOPED_TIMER(_append_data_timer);
698
11.1k
    IndexStreamSharedPtr index_stream;
699
700
11.1k
    int64_t index_id = header.index_id();
701
11.1k
    DBUG_EXECUTE_IF("TabletStream._append_data.unknown_indexid",
702
11.1k
                    { index_id = UNKNOWN_ID_FOR_TEST; });
703
11.1k
    auto it = _index_streams_map.find(index_id);
704
11.1k
    if (it == _index_streams_map.end()) {
705
1
        return Status::Error<ErrorCode::INVALID_ARGUMENT>("unknown index_id {}", index_id);
706
11.1k
    } else {
707
11.1k
        index_stream = it->second;
708
11.1k
    }
709
710
11.1k
    return index_stream->append_data(header, data);
711
11.1k
}
712
713
208
int LoadStream::on_received_messages(StreamId id, butil::IOBuf* const messages[], size_t size) {
714
208
    VLOG_DEBUG << "on_received_messages " << id << " " << size;
715
429
    for (size_t i = 0; i < size; ++i) {
716
11.5k
        while (messages[i]->size() > 0) {
717
            // step 1: parse header
718
11.3k
            size_t hdr_len = 0;
719
11.3k
            messages[i]->cutn((void*)&hdr_len, sizeof(size_t));
720
11.3k
            butil::IOBuf hdr_buf;
721
11.3k
            PStreamHeader hdr;
722
11.3k
            messages[i]->cutn(&hdr_buf, hdr_len);
723
11.3k
            _parse_header(&hdr_buf, hdr);
724
725
            // step 2: cut data
726
11.3k
            size_t data_len = 0;
727
11.3k
            messages[i]->cutn((void*)&data_len, sizeof(size_t));
728
11.3k
            butil::IOBuf data_buf;
729
11.3k
            PStreamHeader data;
730
11.3k
            messages[i]->cutn(&data_buf, data_len);
731
732
            // step 3: dispatch
733
11.3k
            _dispatch(id, hdr, &data_buf);
734
11.3k
        }
735
221
    }
736
208
    return 0;
737
208
}
738
739
11.3k
void LoadStream::_dispatch(StreamId id, const PStreamHeader& hdr, butil::IOBuf* data) {
740
11.3k
    VLOG_DEBUG << PStreamHeader_Opcode_Name(hdr.opcode()) << " from " << hdr.src_id()
741
0
               << " with tablet " << hdr.tablet_id();
742
11.3k
    SCOPED_ATTACH_TASK(_resource_ctx);
743
    // CLOSE_LOAD message should not be fault injected,
744
    // otherwise the message will be ignored and causing close wait timeout
745
11.3k
    if (hdr.opcode() != PStreamHeader::CLOSE_LOAD) {
746
11.1k
        DBUG_EXECUTE_IF("LoadStream._dispatch.unknown_loadid", {
747
11.1k
            PStreamHeader& t_hdr = const_cast<PStreamHeader&>(hdr);
748
11.1k
            PUniqueId* load_id = t_hdr.mutable_load_id();
749
11.1k
            load_id->set_hi(UNKNOWN_ID_FOR_TEST);
750
11.1k
            load_id->set_lo(UNKNOWN_ID_FOR_TEST);
751
11.1k
        });
752
11.1k
        DBUG_EXECUTE_IF("LoadStream._dispatch.unknown_srcid", {
753
11.1k
            PStreamHeader& t_hdr = const_cast<PStreamHeader&>(hdr);
754
11.1k
            t_hdr.set_src_id(UNKNOWN_ID_FOR_TEST);
755
11.1k
        });
756
11.1k
    }
757
11.3k
    if (UniqueId(hdr.load_id()) != UniqueId(_load_id)) {
758
1
        Status st = Status::Error<ErrorCode::INVALID_ARGUMENT>(
759
1
                "invalid load id {}, expected {}", print_id(hdr.load_id()), print_id(_load_id));
760
1
        _report_failure(id, st, hdr);
761
1
        return;
762
1
    }
763
764
11.3k
    {
765
11.3k
        std::lock_guard lock_guard(_lock);
766
11.3k
        if (!_open_streams.contains(hdr.src_id())) {
767
17
            Status st = Status::Error<ErrorCode::INVALID_ARGUMENT>("no open stream from source {}",
768
17
                                                                   hdr.src_id());
769
17
            _report_failure(id, st, hdr);
770
17
            return;
771
17
        }
772
11.3k
    }
773
774
11.2k
    switch (hdr.opcode()) {
775
168
    case PStreamHeader::ADD_SEGMENT: {
776
168
        auto st = _append_data(hdr, data);
777
168
        if (!st.ok()) {
778
0
            _report_failure(id, st, hdr);
779
168
        } else {
780
            // Report tablet load info only on ADD_SEGMENT to reduce frequency.
781
            // ADD_SEGMENT is sent once per segment, while APPEND_DATA is sent
782
            // for every data batch. This reduces unnecessary writes and avoids
783
            // potential stream write failures when the sender is closing.
784
168
            _report_tablet_load_info(id, hdr.index_id());
785
168
        }
786
168
    } break;
787
10.9k
    case PStreamHeader::APPEND_DATA: {
788
10.9k
        auto st = _append_data(hdr, data);
789
10.9k
        if (!st.ok()) {
790
2
            _report_failure(id, st, hdr);
791
2
        }
792
10.9k
    } break;
793
175
    case PStreamHeader::CLOSE_LOAD: {
794
175
        DBUG_EXECUTE_IF("LoadStream.close_load.block", DBUG_BLOCK);
795
175
        std::vector<int64_t> success_tablet_ids;
796
175
        FailedTablets failed_tablets;
797
175
        std::vector<PTabletID> tablets_to_commit(hdr.tablets().begin(), hdr.tablets().end());
798
        // Step 1: count this CLOSE_LOAD and, if this is the last one, commit. Under _lock.
799
175
        bool all_received =
800
175
                close(hdr.src_id(), tablets_to_commit, &success_tablet_ids, &failed_tablets);
801
        // Step 2: send THIS stream's EOS (network IO, must be outside _lock). A stream
802
        // must not be StreamClose'd before its own EOS is delivered, otherwise the
803
        // sender sees on_closed without EOS and reports "Stream closed without EOS".
804
175
        _report_result(id, Status::OK(), success_tablet_ids, failed_tablets, true);
805
175
        bool is_incremental =
806
175
                hdr.has_num_incremental_streams() && hdr.num_incremental_streams() > 0;
807
        // Test-only: delay every incremental stream except the one that made
808
        // all_received, so a non-last incremental stream parks after the last
809
        // stream drained the list. On the buggy code this orphans it and the
810
        // load hangs; on the fix the latch drains the late registration under
811
        // the same lock. Inert unless enable_debug_points=true.
812
175
        if (is_incremental && !all_received) {
813
2
            DBUG_EXECUTE_IF("LoadStream.close_load.delay_incremental_register",
814
2
                            { bthread_usleep(3000000); });
815
2
        }
816
        // Step 3: close this stream (non-incremental) or park it for deferred close
817
        // (incremental, fencing), then collect everything that is now safe to close.
818
        // Registration happens only after step 2, so a collected stream already had
819
        // its EOS delivered (fixes the close-before-EOS race); the all-received latch
820
        // inside makes any late thread drain the parked streams (fixes the leak race).
821
175
        auto streams_to_close = mark_eos_sent_and_collect(id, is_incremental);
822
175
        for (auto& closing_id : streams_to_close) {
823
175
            brpc::StreamClose(closing_id);
824
175
        }
825
175
    } break;
826
0
    case PStreamHeader::GET_SCHEMA: {
827
0
        _report_schema(id, hdr);
828
0
    } break;
829
0
    default:
830
0
        LOG(WARNING) << "unexpected stream message " << hdr.opcode() << ", " << *this;
831
0
        DCHECK(false);
832
11.2k
    }
833
11.2k
}
834
835
0
void LoadStream::on_idle_timeout(StreamId id) {
836
0
    LOG(WARNING) << "closing load stream on idle timeout, " << *this;
837
0
    brpc::StreamClose(id);
838
0
}
839
840
175
void LoadStream::on_closed(StreamId id) {
841
    // `this` may be freed by other threads after increasing `_close_rpc_cnt`,
842
    // format string first to prevent use-after-free
843
175
    std::stringstream ss;
844
175
    ss << *this;
845
175
    auto remaining_streams = _total_streams - _close_rpc_cnt.fetch_add(1) - 1;
846
175
    LOG(INFO) << "stream " << id << " on_closed, remaining streams = " << remaining_streams << ", "
847
175
              << ss.str();
848
175
    if (remaining_streams == 0) {
849
93
        _load_stream_mgr->clear_load(_load_id);
850
93
    }
851
175
}
852
853
824
inline std::ostream& operator<<(std::ostream& ostr, const LoadStream& load_stream) {
854
824
    ostr << "load_id=" << print_id(load_stream._load_id) << ", txn_id=" << load_stream._txn_id;
855
824
    return ostr;
856
824
}
857
858
} // namespace doris