Skip to content

Commit e3fd168

Browse files
api-clients-generation-pipeline[bot]ci.datadog-api-spec
andauthored
Add buffer configuration to ClickHouse destination (#1783)
Co-authored-by: ci.datadog-api-spec <packages@datadoghq.com>
1 parent debc7d5 commit e3fd168

6 files changed

Lines changed: 228 additions & 0 deletions

.generator/schemas/v2/openapi.yaml

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64398,6 +64398,8 @@ components:
6439864398
$ref: "#/components/schemas/ObservabilityPipelineClickhouseDestinationBatch"
6439964399
batch_encoding:
6440064400
$ref: "#/components/schemas/ObservabilityPipelineClickhouseDestinationBatchEncoding"
64401+
buffer:
64402+
$ref: "#/components/schemas/ObservabilityPipelineBufferOptions"
6440164403
compression:
6440264404
$ref: "#/components/schemas/ObservabilityPipelineClickhouseDestinationCompression"
6440364405
database:
Lines changed: 148 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,148 @@
1+
// Validate an observability pipeline with ClickHouse destination with all fields
2+
// set returns "OK" response
3+
use datadog_api_client::datadog;
4+
use datadog_api_client::datadogV2::api_observability_pipelines::ObservabilityPipelinesAPI;
5+
use datadog_api_client::datadogV2::model::ObservabilityPipelineBufferOptions;
6+
use datadog_api_client::datadogV2::model::ObservabilityPipelineBufferOptionsMemoryType;
7+
use datadog_api_client::datadogV2::model::ObservabilityPipelineBufferOptionsWhenFull;
8+
use datadog_api_client::datadogV2::model::ObservabilityPipelineClickhouseDestination;
9+
use datadog_api_client::datadogV2::model::ObservabilityPipelineClickhouseDestinationAuth;
10+
use datadog_api_client::datadogV2::model::ObservabilityPipelineClickhouseDestinationAuthStrategy;
11+
use datadog_api_client::datadogV2::model::ObservabilityPipelineClickhouseDestinationBatch;
12+
use datadog_api_client::datadogV2::model::ObservabilityPipelineClickhouseDestinationBatchEncoding;
13+
use datadog_api_client::datadogV2::model::ObservabilityPipelineClickhouseDestinationBatchEncodingCodec;
14+
use datadog_api_client::datadogV2::model::ObservabilityPipelineClickhouseDestinationCompression;
15+
use datadog_api_client::datadogV2::model::ObservabilityPipelineClickhouseDestinationCompressionAlgorithm;
16+
use datadog_api_client::datadogV2::model::ObservabilityPipelineClickhouseDestinationCompressionObject;
17+
use datadog_api_client::datadogV2::model::ObservabilityPipelineClickhouseDestinationFormat;
18+
use datadog_api_client::datadogV2::model::ObservabilityPipelineClickhouseDestinationType;
19+
use datadog_api_client::datadogV2::model::ObservabilityPipelineConfig;
20+
use datadog_api_client::datadogV2::model::ObservabilityPipelineConfigDestinationItem;
21+
use datadog_api_client::datadogV2::model::ObservabilityPipelineConfigProcessorGroup;
22+
use datadog_api_client::datadogV2::model::ObservabilityPipelineConfigProcessorItem;
23+
use datadog_api_client::datadogV2::model::ObservabilityPipelineConfigSourceItem;
24+
use datadog_api_client::datadogV2::model::ObservabilityPipelineDataAttributes;
25+
use datadog_api_client::datadogV2::model::ObservabilityPipelineDatadogAgentSource;
26+
use datadog_api_client::datadogV2::model::ObservabilityPipelineDatadogAgentSourceType;
27+
use datadog_api_client::datadogV2::model::ObservabilityPipelineFilterProcessor;
28+
use datadog_api_client::datadogV2::model::ObservabilityPipelineFilterProcessorType;
29+
use datadog_api_client::datadogV2::model::ObservabilityPipelineMemoryBufferSizeOptions;
30+
use datadog_api_client::datadogV2::model::ObservabilityPipelineSpec;
31+
use datadog_api_client::datadogV2::model::ObservabilityPipelineSpecData;
32+
use datadog_api_client::datadogV2::model::ObservabilityPipelineTls;
33+
34+
#[tokio::main]
35+
async fn main() {
36+
let body =
37+
ObservabilityPipelineSpec::new(
38+
ObservabilityPipelineSpecData::new(
39+
ObservabilityPipelineDataAttributes::new(
40+
ObservabilityPipelineConfig::new(
41+
vec![
42+
ObservabilityPipelineConfigDestinationItem::ObservabilityPipelineClickhouseDestination(
43+
Box::new(
44+
ObservabilityPipelineClickhouseDestination::new(
45+
"clickhouse-destination".to_string(),
46+
vec!["my-processor-group".to_string()],
47+
"application_logs".to_string(),
48+
ObservabilityPipelineClickhouseDestinationType::CLICKHOUSE,
49+
)
50+
.auth(
51+
ObservabilityPipelineClickhouseDestinationAuth::new(
52+
ObservabilityPipelineClickhouseDestinationAuthStrategy::BASIC,
53+
)
54+
.password_key("CLICKHOUSE_PASSWORD".to_string())
55+
.username_key("CLICKHOUSE_USERNAME".to_string()),
56+
)
57+
.batch(
58+
ObservabilityPipelineClickhouseDestinationBatch::new()
59+
.max_events(1000)
60+
.timeout_secs(1),
61+
)
62+
.batch_encoding(
63+
ObservabilityPipelineClickhouseDestinationBatchEncoding::new(
64+
ObservabilityPipelineClickhouseDestinationBatchEncodingCodec
65+
::ARROW_STREAM,
66+
).allow_nullable_fields(true),
67+
)
68+
.buffer(
69+
ObservabilityPipelineBufferOptions
70+
::ObservabilityPipelineMemoryBufferSizeOptions(
71+
Box::new(
72+
ObservabilityPipelineMemoryBufferSizeOptions::new(500)
73+
.type_(ObservabilityPipelineBufferOptionsMemoryType::MEMORY)
74+
.when_full(ObservabilityPipelineBufferOptionsWhenFull::BLOCK),
75+
),
76+
),
77+
)
78+
.compression(
79+
ObservabilityPipelineClickhouseDestinationCompression
80+
::ObservabilityPipelineClickhouseDestinationCompressionObject(
81+
Box::new(
82+
ObservabilityPipelineClickhouseDestinationCompressionObject::new(
83+
ObservabilityPipelineClickhouseDestinationCompressionAlgorithm
84+
::GZIP,
85+
).level(6),
86+
),
87+
),
88+
)
89+
.database("my_database".to_string())
90+
.date_time_best_effort(true)
91+
.endpoint_url_key("CLICKHOUSE_ENDPOINT_URL".to_string())
92+
.format(ObservabilityPipelineClickhouseDestinationFormat::ARROW_STREAM)
93+
.skip_unknown_fields(Some(true))
94+
.tls(
95+
ObservabilityPipelineTls::new("/path/to/cert.crt".to_string())
96+
.ca_file("/path/to/ca.crt".to_string())
97+
.key_file("/path/to/key.key".to_string())
98+
.key_pass_key("TLS_KEY_PASSPHRASE".to_string()),
99+
),
100+
),
101+
)
102+
],
103+
vec![
104+
ObservabilityPipelineConfigSourceItem::ObservabilityPipelineDatadogAgentSource(
105+
Box::new(
106+
ObservabilityPipelineDatadogAgentSource::new(
107+
"datadog-agent-source".to_string(),
108+
ObservabilityPipelineDatadogAgentSourceType::DATADOG_AGENT,
109+
),
110+
),
111+
)
112+
],
113+
).processor_groups(
114+
vec![
115+
ObservabilityPipelineConfigProcessorGroup::new(
116+
true,
117+
"my-processor-group".to_string(),
118+
"service:my-service".to_string(),
119+
vec!["datadog-agent-source".to_string()],
120+
vec![
121+
ObservabilityPipelineConfigProcessorItem::ObservabilityPipelineFilterProcessor(
122+
Box::new(
123+
ObservabilityPipelineFilterProcessor::new(
124+
true,
125+
"filter-processor".to_string(),
126+
"status:error".to_string(),
127+
ObservabilityPipelineFilterProcessorType::FILTER,
128+
),
129+
),
130+
)
131+
],
132+
)
133+
],
134+
),
135+
"Pipeline with ClickHouse Destination All Fields".to_string(),
136+
),
137+
"pipelines".to_string(),
138+
),
139+
);
140+
let configuration = datadog::Configuration::new();
141+
let api = ObservabilityPipelinesAPI::with_config(configuration);
142+
let resp = api.validate_pipeline(body).await;
143+
if let Ok(value) = resp {
144+
println!("{:#?}", value);
145+
} else {
146+
println!("{:#?}", resp.unwrap_err());
147+
}
148+
}

src/datadogV2/model/model_observability_pipeline_clickhouse_destination.rs

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,9 @@ pub struct ObservabilityPipelineClickhouseDestination {
2525
#[serde(rename = "batch_encoding")]
2626
pub batch_encoding:
2727
Option<crate::datadogV2::model::ObservabilityPipelineClickhouseDestinationBatchEncoding>,
28+
/// Configuration for buffer settings on destination components.
29+
#[serde(rename = "buffer")]
30+
pub buffer: Option<crate::datadogV2::model::ObservabilityPipelineBufferOptions>,
2831
/// Compression setting for outbound HTTP requests to ClickHouse.
2932
/// Can be specified as a shorthand string (`"gzip"` or `"none"`) or as an object
3033
/// with an `algorithm` field and an optional `level` (gzip only, 1–9).
@@ -89,6 +92,7 @@ impl ObservabilityPipelineClickhouseDestination {
8992
auth: None,
9093
batch: None,
9194
batch_encoding: None,
95+
buffer: None,
9296
compression: None,
9397
database: None,
9498
date_time_best_effort: None,
@@ -129,6 +133,14 @@ impl ObservabilityPipelineClickhouseDestination {
129133
self
130134
}
131135

136+
pub fn buffer(
137+
mut self,
138+
value: crate::datadogV2::model::ObservabilityPipelineBufferOptions,
139+
) -> Self {
140+
self.buffer = Some(value);
141+
self
142+
}
143+
132144
pub fn compression(
133145
mut self,
134146
value: crate::datadogV2::model::ObservabilityPipelineClickhouseDestinationCompression,
@@ -203,6 +215,9 @@ impl<'de> Deserialize<'de> for ObservabilityPipelineClickhouseDestination {
203215
crate::datadogV2::model::ObservabilityPipelineClickhouseDestinationBatch,
204216
> = None;
205217
let mut batch_encoding: Option<crate::datadogV2::model::ObservabilityPipelineClickhouseDestinationBatchEncoding> = None;
218+
let mut buffer: Option<
219+
crate::datadogV2::model::ObservabilityPipelineBufferOptions,
220+
> = None;
206221
let mut compression: Option<
207222
crate::datadogV2::model::ObservabilityPipelineClickhouseDestinationCompression,
208223
> = None;
@@ -247,6 +262,20 @@ impl<'de> Deserialize<'de> for ObservabilityPipelineClickhouseDestination {
247262
batch_encoding =
248263
Some(serde_json::from_value(v).map_err(M::Error::custom)?);
249264
}
265+
"buffer" => {
266+
if v.is_null() {
267+
continue;
268+
}
269+
buffer = Some(serde_json::from_value(v).map_err(M::Error::custom)?);
270+
if let Some(ref _buffer) = buffer {
271+
match _buffer {
272+
crate::datadogV2::model::ObservabilityPipelineBufferOptions::UnparsedObject(_buffer) => {
273+
_unparsed = true;
274+
},
275+
_ => {}
276+
}
277+
}
278+
}
250279
"compression" => {
251280
if v.is_null() {
252281
continue;
@@ -342,6 +371,7 @@ impl<'de> Deserialize<'de> for ObservabilityPipelineClickhouseDestination {
342371
auth,
343372
batch,
344373
batch_encoding,
374+
buffer,
345375
compression,
346376
database,
347377
date_time_best_effort,
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
2026-06-24T16:45:05.037Z
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
1+
{
2+
"http_interactions": [
3+
{
4+
"request": {
5+
"body": {
6+
"string": "{\"data\":{\"attributes\":{\"config\":{\"destinations\":[{\"auth\":{\"password_key\":\"CLICKHOUSE_PASSWORD\",\"strategy\":\"basic\",\"username_key\":\"CLICKHOUSE_USERNAME\"},\"batch\":{\"max_events\":1000,\"timeout_secs\":1},\"batch_encoding\":{\"allow_nullable_fields\":true,\"codec\":\"arrow_stream\"},\"buffer\":{\"max_events\":500,\"type\":\"memory\",\"when_full\":\"block\"},\"compression\":{\"algorithm\":\"gzip\",\"level\":6},\"database\":\"my_database\",\"date_time_best_effort\":true,\"endpoint_url_key\":\"CLICKHOUSE_ENDPOINT_URL\",\"format\":\"arrow_stream\",\"id\":\"clickhouse-destination\",\"inputs\":[\"my-processor-group\"],\"skip_unknown_fields\":true,\"table\":\"application_logs\",\"tls\":{\"ca_file\":\"/path/to/ca.crt\",\"crt_file\":\"/path/to/cert.crt\",\"key_file\":\"/path/to/key.key\",\"key_pass_key\":\"TLS_KEY_PASSPHRASE\"},\"type\":\"clickhouse\"}],\"processor_groups\":[{\"enabled\":true,\"id\":\"my-processor-group\",\"include\":\"service:my-service\",\"inputs\":[\"datadog-agent-source\"],\"processors\":[{\"enabled\":true,\"id\":\"filter-processor\",\"include\":\"status:error\",\"type\":\"filter\"}]}],\"sources\":[{\"id\":\"datadog-agent-source\",\"type\":\"datadog_agent\"}]},\"name\":\"Pipeline with ClickHouse Destination All Fields\"},\"type\":\"pipelines\"}}",
7+
"encoding": null
8+
},
9+
"headers": {
10+
"Accept": [
11+
"application/json"
12+
],
13+
"Content-Type": [
14+
"application/json"
15+
]
16+
},
17+
"method": "post",
18+
"uri": "https://api.datadoghq.com/api/v2/obs-pipelines/pipelines/validate"
19+
},
20+
"response": {
21+
"body": {
22+
"string": "{\"errors\":[]}\n",
23+
"encoding": null
24+
},
25+
"headers": {
26+
"Content-Type": [
27+
"application/vnd.api+json"
28+
]
29+
},
30+
"status": {
31+
"code": 200,
32+
"message": "OK"
33+
}
34+
},
35+
"recorded_at": "Wed, 24 Jun 2026 16:45:05 GMT"
36+
}
37+
],
38+
"recorded_with": "VCR 6.0.0"
39+
}

tests/scenarios/features/v2/observability_pipelines.feature

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -192,6 +192,14 @@ Feature: Observability Pipelines
192192
Then the response status is 200 OK
193193
And the response "errors" has length 0
194194

195+
@team:DataDog/observability-pipelines
196+
Scenario: Validate an observability pipeline with ClickHouse destination with all fields set returns "OK" response
197+
Given new "ValidatePipeline" request
198+
And body with value {"data": {"attributes": {"config": {"destinations": [{"id": "clickhouse-destination", "inputs": ["my-processor-group"], "type": "clickhouse", "endpoint_url_key": "CLICKHOUSE_ENDPOINT_URL", "database": "my_database", "table": "application_logs", "format": "arrow_stream", "skip_unknown_fields": true, "date_time_best_effort": true, "compression": {"algorithm": "gzip", "level": 6}, "auth": {"strategy": "basic", "username_key": "CLICKHOUSE_USERNAME", "password_key": "CLICKHOUSE_PASSWORD"}, "batch": {"max_events": 1000, "timeout_secs": 1}, "batch_encoding": {"codec": "arrow_stream", "allow_nullable_fields": true}, "tls": {"crt_file": "/path/to/cert.crt", "ca_file": "/path/to/ca.crt", "key_file": "/path/to/key.key", "key_pass_key": "TLS_KEY_PASSPHRASE"}, "buffer": {"type": "memory", "max_events": 500, "when_full": "block"}}], "processor_groups": [{"enabled": true, "id": "my-processor-group", "include": "service:my-service", "inputs": ["datadog-agent-source"], "processors": [{"enabled": true, "id": "filter-processor", "include": "status:error", "type": "filter"}]}], "sources": [{"id": "datadog-agent-source", "type": "datadog_agent"}]}, "name": "Pipeline with ClickHouse Destination All Fields"}, "type": "pipelines"}}
199+
When the request is sent
200+
Then the response status is 200 OK
201+
And the response "errors" has length 0
202+
195203
@team:DataDog/observability-pipelines
196204
Scenario: Validate an observability pipeline with HTTP server source valid_tokens returns "OK" response
197205
Given new "ValidatePipeline" request

0 commit comments

Comments
 (0)