Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions bazel/container_images.bzl
Original file line number Diff line number Diff line change
Expand Up @@ -367,3 +367,11 @@ def stirling_test_images():
repository = "golang_1_22_grpc_server_with_buildinfo",
digest = "sha256:67adba5e8513670fa37bd042862e7844f26239e8d2997ed8c3b0aa527bc04cc3",
)

# ClickHouse server image for testing.
# clickhouse/clickhouse-server:25.7-alpine
_container_image(
name = "clickhouse_server_image_25_7",
repository = "clickhouse_server_image_25_7",
digest = "sha256:60c53a520a1caad6555eb6772a8a9c91bb09774c1c7ec87e3371ea3da254eeab",
)
2 changes: 2 additions & 0 deletions src/carnot/carnot.cc
Original file line number Diff line number Diff line change
Expand Up @@ -181,6 +181,8 @@ Status CarnotImpl::RegisterUDFsInPlanFragment(exec::ExecState* exec_state, plan:
.OnUDTFSource(no_op)
.OnEmptySource(no_op)
.OnOTelSink(no_op)
.OnClickHouseSource(no_op)
.OnClickHouseExportSink(no_op)
.Walk(pf);
}

Expand Down
44 changes: 44 additions & 0 deletions src/carnot/exec/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@ pl_cc_library(
"//src/shared/types:cc_library",
"//src/table_store/table:cc_library",
"@com_github_apache_arrow//:arrow",
"@com_github_clickhouse_clickhouse_cpp//:clickhouse_cpp",
"@com_github_grpc_grpc//:grpc++",
"@com_github_opentelemetry_proto//:logs_service_grpc_cc",
"@com_github_opentelemetry_proto//:metrics_service_grpc_cc",
Expand Down Expand Up @@ -300,3 +301,46 @@ pl_cc_test(
"@com_github_grpc_grpc//:grpc++_test",
],
)

pl_cc_test(
name = "clickhouse_source_node_test",
timeout = "long",
srcs = ["clickhouse_source_node_test.cc"],
data = [
"//src/stirling/source_connectors/socket_tracer/testing/container_images/clickhouse",
],
tags = [
"exclusive",
"requires_bpf",
],
deps = [
":cc_library",
":exec_node_test_helpers",
":test_utils",
"//src/carnot/planpb:plan_testutils",
"//src/common/testing/test_utils:cc_library",
"@com_github_clickhouse_clickhouse_cpp//:clickhouse_cpp",
],
)

pl_cc_test(
name = "clickhouse_export_sink_node_test",
timeout = "long",
srcs = ["clickhouse_export_sink_node_test.cc"],
data = [
"//src/stirling/source_connectors/socket_tracer/testing/container_images/clickhouse",
],
tags = [
"exclusive",
"requires_bpf",
],
deps = [
":cc_library",
":exec_node_test_helpers",
":test_utils",
"//src/carnot/plan:cc_library",
"//src/carnot/planpb:plan_pl_cc_proto",
"//src/common/testing/test_utils:cc_library",
"@com_github_clickhouse_clickhouse_cpp//:clickhouse_cpp",
],
)
10 changes: 10 additions & 0 deletions src/carnot/exec/exec_graph.cc
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,8 @@
#include <unordered_map>

#include "src/carnot/exec/agg_node.h"
#include "src/carnot/exec/clickhouse_export_sink_node.h"
#include "src/carnot/exec/clickhouse_source_node.h"
#include "src/carnot/exec/empty_source_node.h"
#include "src/carnot/exec/equijoin_node.h"
#include "src/carnot/exec/exec_node.h"
Expand Down Expand Up @@ -108,6 +110,14 @@ Status ExecutionGraph::Init(table_store::schema::Schema* schema, plan::PlanState
.OnOTelSink([&](auto& node) {
return OnOperatorImpl<plan::OTelExportSinkOperator, OTelExportSinkNode>(node, &descriptors);
})
.OnClickHouseSource([&](auto& node) {
return OnOperatorImpl<plan::ClickHouseSourceOperator, ClickHouseSourceNode>(node,
&descriptors);
})
.OnClickHouseExportSink([&](auto& node) {
return OnOperatorImpl<plan::ClickHouseExportSinkOperator, ClickHouseExportSinkNode>(node,
&descriptors);
})
.Walk(pf_);
}

Expand Down
65 changes: 65 additions & 0 deletions src/carnot/plan/operators.cc
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,10 @@ std::unique_ptr<Operator> Operator::FromProto(const planpb::Operator& pb, int64_
return CreateOperator<UDTFSourceOperator>(id, pb.udtf_source_op());
case planpb::EMPTY_SOURCE_OPERATOR:
return CreateOperator<EmptySourceOperator>(id, pb.empty_source_op());
case planpb::CLICKHOUSE_SOURCE_OPERATOR:
return CreateOperator<ClickHouseSourceOperator>(id, pb.clickhouse_source_op());
case planpb::CLICKHOUSE_EXPORT_SINK_OPERATOR:
return CreateOperator<ClickHouseExportSinkOperator>(id, pb.clickhouse_sink_op());
case planpb::OTEL_EXPORT_SINK_OPERATOR:
return CreateOperator<OTelExportSinkOperator>(id, pb.otel_sink_op());
default:
Expand Down Expand Up @@ -709,6 +713,67 @@ StatusOr<table_store::schema::Relation> EmptySourceOperator::OutputRelation(
return r;
}

/**
* ClickHouseSourceOperator implementation.
*/

std::string ClickHouseSourceOperator::DebugString() const {
return absl::Substitute(R"(Op:ClickHouseSource(
host=$0
port=$1
username=$2
batch_size=$3
start_time=$4
end_time=$5
timestamp_column=$6
partition_column=$7
)",
pb_.host(), pb_.port(), pb_.username(), pb_.batch_size(),
pb_.start_time(), pb_.end_time(), pb_.timestamp_column(),
pb_.partition_column());
}

Status ClickHouseSourceOperator::Init(const planpb::ClickHouseSourceOperator& pb) {
pb_ = pb;
is_initialized_ = true;
return Status::OK();
}

StatusOr<table_store::schema::Relation> ClickHouseSourceOperator::OutputRelation(
const table_store::schema::Schema&, const PlanState&,
const std::vector<int64_t>& input_ids) const {
DCHECK(is_initialized_) << "Not initialized";
if (!input_ids.empty()) {
return error::InvalidArgument("Source operator cannot have any inputs");
}
table_store::schema::Relation r;
for (int i = 0; i < pb_.column_types_size(); ++i) {
r.AddColumn(static_cast<types::DataType>(pb_.column_types(i)), pb_.column_names(i));
}
return r;
}

/**
* ClickHouse Export Sink Operator Implementation.
*/

Status ClickHouseExportSinkOperator::Init(const planpb::ClickHouseExportSinkOperator& pb) {
pb_ = pb;
is_initialized_ = true;
return Status::OK();
}

StatusOr<table_store::schema::Relation> ClickHouseExportSinkOperator::OutputRelation(
const table_store::schema::Schema&, const PlanState&, const std::vector<int64_t>&) const {
DCHECK(is_initialized_) << "Not initialized";
// There are no outputs.
return table_store::schema::Relation();
}

std::string ClickHouseExportSinkOperator::DebugString() const {
return absl::Substitute("Op:ClickHouseExportSink(table=$0)", pb_.table_name());
}

/**
* OTel Export Sink Operator Implementation.
*/
Expand Down
63 changes: 63 additions & 0 deletions src/carnot/plan/operators.h
Original file line number Diff line number Diff line change
Expand Up @@ -359,6 +359,69 @@ class EmptySourceOperator : public Operator {
std::vector<int64_t> column_idxs_;
};

class ClickHouseSourceOperator : public Operator {
public:
explicit ClickHouseSourceOperator(int64_t id)
: Operator(id, planpb::CLICKHOUSE_SOURCE_OPERATOR) {}
~ClickHouseSourceOperator() override = default;

StatusOr<table_store::schema::Relation> OutputRelation(
const table_store::schema::Schema& schema, const PlanState& state,
const std::vector<int64_t>& input_ids) const override;
Status Init(const planpb::ClickHouseSourceOperator& pb);
std::string DebugString() const override;

std::string host() const { return pb_.host(); }
int32_t port() const { return pb_.port(); }
std::string username() const { return pb_.username(); }
std::string password() const { return pb_.password(); }
std::string database() const { return pb_.database(); }
std::string query() const { return pb_.query(); }
int32_t batch_size() const { return pb_.batch_size(); }
bool streaming() const { return pb_.streaming(); }
std::vector<std::string> column_names() const {
return std::vector<std::string>(pb_.column_names().begin(), pb_.column_names().end());
}
std::vector<types::DataType> column_types() const {
std::vector<types::DataType> types;
types.reserve(pb_.column_types_size());
for (const auto& type : pb_.column_types()) {
types.push_back(static_cast<types::DataType>(type));
}
return types;
}
std::string timestamp_column() const { return pb_.timestamp_column(); }
std::string partition_column() const { return pb_.partition_column(); }
int64_t start_time() const { return pb_.start_time(); }
int64_t end_time() const { return pb_.end_time(); }

private:
planpb::ClickHouseSourceOperator pb_;
};

class ClickHouseExportSinkOperator : public Operator {
public:
explicit ClickHouseExportSinkOperator(int64_t id)
: Operator(id, planpb::CLICKHOUSE_EXPORT_SINK_OPERATOR) {}
~ClickHouseExportSinkOperator() override = default;

StatusOr<table_store::schema::Relation> OutputRelation(
const table_store::schema::Schema& schema, const PlanState& state,
const std::vector<int64_t>& input_ids) const override;
Status Init(const planpb::ClickHouseExportSinkOperator& pb);
std::string DebugString() const override;

const planpb::ClickHouseConfig& clickhouse_config() const { return pb_.clickhouse_config(); }
const std::string& table_name() const { return pb_.table_name(); }
const ::google::protobuf::RepeatedPtrField<planpb::ClickHouseExportSinkOperator::ColumnMapping>&
column_mappings() const {
return pb_.column_mappings();
}

private:
planpb::ClickHouseExportSinkOperator pb_;
};

class OTelExportSinkOperator : public Operator {
public:
explicit OTelExportSinkOperator(int64_t id) : Operator(id, planpb::OTEL_EXPORT_SINK_OPERATOR) {}
Expand Down
7 changes: 7 additions & 0 deletions src/carnot/plan/plan_fragment.cc
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,13 @@ Status PlanFragmentWalker::CallWalkFn(const Operator& op) {
case planpb::OperatorType::OTEL_EXPORT_SINK_OPERATOR:
PX_RETURN_IF_ERROR(CallAs<OTelExportSinkOperator>(on_otel_sink_walk_fn_, op));
break;
case planpb::OperatorType::CLICKHOUSE_SOURCE_OPERATOR:
PX_RETURN_IF_ERROR(CallAs<ClickHouseSourceOperator>(on_clickhouse_source_walk_fn_, op));
break;
case planpb::OperatorType::CLICKHOUSE_EXPORT_SINK_OPERATOR:
PX_RETURN_IF_ERROR(
CallAs<ClickHouseExportSinkOperator>(on_clickhouse_export_sink_walk_fn_, op));
break;
default:
LOG(FATAL) << absl::Substitute("Operator does not exist: $0", magic_enum::enum_name(op_type));
return error::InvalidArgument("Operator does not exist: $0", magic_enum::enum_name(op_type));
Expand Down
15 changes: 15 additions & 0 deletions src/carnot/plan/plan_fragment.h
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,8 @@ class PlanFragmentWalker {
using UDTFSourceWalkFn = std::function<Status(const UDTFSourceOperator&)>;
using EmptySourceWalkFn = std::function<Status(const EmptySourceOperator&)>;
using OTelSinkWalkFn = std::function<Status(const OTelExportSinkOperator&)>;
using ClickHouseSourceWalkFn = std::function<Status(const ClickHouseSourceOperator&)>;
using ClickHouseExportSinkWalkFn = std::function<Status(const ClickHouseExportSinkOperator&)>;

/**
* Register callback for when a memory source operator is encountered.
Expand Down Expand Up @@ -181,6 +183,17 @@ class PlanFragmentWalker {
on_otel_sink_walk_fn_ = fn;
return *this;
}

PlanFragmentWalker& OnClickHouseSource(const ClickHouseSourceWalkFn& fn) {
on_clickhouse_source_walk_fn_ = fn;
return *this;
}

PlanFragmentWalker& OnClickHouseExportSink(const ClickHouseExportSinkWalkFn& fn) {
on_clickhouse_export_sink_walk_fn_ = fn;
return *this;
}

/**
* Perform a walk of the plan fragment operators in a topologically-sorted order.
* @param plan_fragment The plan fragment to walk.
Expand All @@ -206,6 +219,8 @@ class PlanFragmentWalker {
UDTFSourceWalkFn on_udtf_source_walk_fn_;
EmptySourceWalkFn on_empty_source_walk_fn_;
OTelSinkWalkFn on_otel_sink_walk_fn_;
ClickHouseSourceWalkFn on_clickhouse_source_walk_fn_;
ClickHouseExportSinkWalkFn on_clickhouse_export_sink_walk_fn_;
};

} // namespace plan
Expand Down
8 changes: 6 additions & 2 deletions src/carnot/planner/compiler_state/compiler_state.h
Original file line number Diff line number Diff line change
Expand Up @@ -119,7 +119,8 @@ class CompilerState : public NotCopyable {
int64_t max_output_rows_per_table, std::string_view result_address,
std::string_view result_ssl_targetname, const RedactionOptions& redaction_options,
std::unique_ptr<planpb::OTelEndpointConfig> endpoint_config,
std::unique_ptr<PluginConfig> plugin_config, DebugInfo debug_info)
std::unique_ptr<PluginConfig> plugin_config, DebugInfo debug_info,
std::unique_ptr<planpb::ClickHouseConfig> clickhouse_config = nullptr)
: relation_map_(std::move(relation_map)),
table_names_to_sensitive_columns_(table_names_to_sensitive_columns),
registry_info_(registry_info),
Expand All @@ -130,7 +131,8 @@ class CompilerState : public NotCopyable {
redaction_options_(redaction_options),
endpoint_config_(std::move(endpoint_config)),
plugin_config_(std::move(plugin_config)),
debug_info_(std::move(debug_info)) {}
debug_info_(std::move(debug_info)),
clickhouse_config_(std::move(clickhouse_config)) {}

CompilerState() = delete;

Expand Down Expand Up @@ -175,6 +177,7 @@ class CompilerState : public NotCopyable {
planpb::OTelEndpointConfig* endpoint_config() { return endpoint_config_.get(); }
PluginConfig* plugin_config() { return plugin_config_.get(); }
const DebugInfo& debug_info() { return debug_info_; }
planpb::ClickHouseConfig* clickhouse_config() { return clickhouse_config_.get(); }

private:
std::unique_ptr<RelationMap> relation_map_;
Expand All @@ -191,6 +194,7 @@ class CompilerState : public NotCopyable {
std::unique_ptr<planpb::OTelEndpointConfig> endpoint_config_ = nullptr;
std::unique_ptr<PluginConfig> plugin_config_ = nullptr;
DebugInfo debug_info_;
std::unique_ptr<planpb::ClickHouseConfig> clickhouse_config_ = nullptr;
};

} // namespace planner
Expand Down
Loading
Loading