aboutsummaryrefslogtreecommitdiffstats
path: root/ydb/library/yql/dq/runtime/dq_input_channel.h
blob: 27020103ccd6f5e7139003218f5298187b041626 (plain) (blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
#pragma once
#include "dq_input.h"
#include "dq_transport.h"

#include <ydb/library/yql/minikql/computation/mkql_computation_node_holders.h>
#include <ydb/library/yql/minikql/mkql_node.h>
#include <ydb/library/yql/utils/yql_panic.h>

namespace NYql::NDq {

struct TDqInputChannelStats : TDqInputStats {
    ui64 ChannelId = 0;

    // profile stats
    TDuration DeserializationTime;

    explicit TDqInputChannelStats(ui64 channelId)
        : ChannelId(channelId) {}
};

class IDqInputChannel : public IDqInput {
public:
    using TPtr = TIntrusivePtr<IDqInputChannel>;

    virtual ui64 GetChannelId() const = 0;

    virtual void Push(NDqProto::TData&& data) = 0;

    virtual void Finish() = 0;

    virtual const TDqInputChannelStats* GetStats() const = 0;
};

IDqInputChannel::TPtr CreateDqInputChannel(ui64 channelId, NKikimr::NMiniKQL::TType* inputType, ui64 maxBufferBytes,
    bool collectProfileStats, const NKikimr::NMiniKQL::TTypeEnvironment& typeEnv,
    const NKikimr::NMiniKQL::THolderFactory& holderFactory, NDqProto::EDataTransportVersion transportVersion);

} // namespace NYql::NDq