Coverage Report

Created: 2026-03-13 19:41

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/io/fs/kafka_consumer_pipe.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 "io/fs/stream_load_pipe.h"
21
22
namespace doris {
23
namespace io {
24
class KafkaConsumerPipe : public StreamLoadPipe {
25
public:
26
    KafkaConsumerPipe(size_t max_buffered_bytes = 1024 * 1024, size_t min_chunk_size = 64 * 1024)
27
5
            : StreamLoadPipe(max_buffered_bytes, min_chunk_size) {}
28
29
5
    ~KafkaConsumerPipe() override = default;
30
31
0
    virtual Status append_with_line_delimiter(const char* data, size_t size) {
32
0
        Status st = append(data, size);
33
0
        if (!st.ok()) {
34
0
            return st;
35
0
        }
36
37
        // append the line delimiter
38
0
        st = append("\n", 1);
39
0
        return st;
40
0
    }
41
42
8
    virtual Status append_json(const char* data, size_t size) {
43
8
        return append_and_flush(data, size);
44
8
    }
45
};
46
} // namespace io
47
} // end namespace doris