Skip to main content

laminar_connectors/otel/schema/
mod.rs

1//! Arrow schemas for OTel signal types.
2//!
3//! Each schema flattens the nested OTLP protobuf hierarchy
4//! (Resource → Scope → Span/DataPoint/LogRecord) into flat columns.
5
6use std::sync::Arc;
7
8use arrow_schema::{DataType, Field, Schema, SchemaRef, TimeUnit};
9
10/// Trace span schema.
11#[must_use]
12pub fn traces_schema() -> SchemaRef {
13    Arc::new(Schema::new(vec![
14        Field::new("trace_id", DataType::FixedSizeBinary(16), false),
15        Field::new("span_id", DataType::FixedSizeBinary(8), false),
16        Field::new("parent_span_id", DataType::FixedSizeBinary(8), true),
17        Field::new("trace_state", DataType::Utf8, true),
18        Field::new("name", DataType::Utf8, false),
19        Field::new("kind", DataType::Int32, false),
20        Field::new("start_time_unix_nano", DataType::Int64, false),
21        Field::new("end_time_unix_nano", DataType::Int64, false),
22        Field::new("duration_ns", DataType::Int64, false),
23        Field::new("status_code", DataType::Int32, false),
24        Field::new("status_message", DataType::Utf8, true),
25        Field::new("resource_service_name", DataType::Utf8, true),
26        Field::new("resource_service_version", DataType::Utf8, true),
27        Field::new("resource_attributes", DataType::Utf8, true),
28        Field::new("scope_name", DataType::Utf8, true),
29        Field::new("scope_version", DataType::Utf8, true),
30        Field::new("attributes", DataType::Utf8, true),
31        Field::new("events_count", DataType::Int32, false),
32        Field::new("links_count", DataType::Int32, false),
33        Field::new(
34            "_laminar_received_at",
35            DataType::Timestamp(TimeUnit::Nanosecond, None),
36            false,
37        ),
38    ]))
39}
40
41/// Metric data point schema.
42#[must_use]
43pub fn metrics_schema() -> SchemaRef {
44    Arc::new(Schema::new(vec![
45        Field::new("metric_name", DataType::Utf8, false),
46        Field::new("metric_description", DataType::Utf8, true),
47        Field::new("metric_unit", DataType::Utf8, true),
48        Field::new("metric_type", DataType::Int32, false),
49        Field::new("timestamp_unix_nano", DataType::Int64, false),
50        Field::new("value_double", DataType::Float64, true),
51        Field::new("value_int", DataType::Int64, true),
52        Field::new("histogram_count", DataType::UInt64, true),
53        Field::new("histogram_sum", DataType::Float64, true),
54        Field::new("resource_service_name", DataType::Utf8, true),
55        Field::new("resource_attributes", DataType::Utf8, true),
56        Field::new("scope_name", DataType::Utf8, true),
57        Field::new("attributes", DataType::Utf8, true),
58        Field::new(
59            "_laminar_received_at",
60            DataType::Timestamp(TimeUnit::Nanosecond, None),
61            false,
62        ),
63    ]))
64}
65
66/// Log record schema.
67#[must_use]
68pub fn logs_schema() -> SchemaRef {
69    Arc::new(Schema::new(vec![
70        Field::new("timestamp_unix_nano", DataType::Int64, false),
71        Field::new("observed_timestamp_unix_nano", DataType::Int64, true),
72        Field::new("severity_number", DataType::Int32, false),
73        Field::new("severity_text", DataType::Utf8, true),
74        Field::new("body_string", DataType::Utf8, true),
75        Field::new("trace_id", DataType::FixedSizeBinary(16), true),
76        Field::new("span_id", DataType::FixedSizeBinary(8), true),
77        Field::new("resource_service_name", DataType::Utf8, true),
78        Field::new("resource_attributes", DataType::Utf8, true),
79        Field::new("scope_name", DataType::Utf8, true),
80        Field::new("attributes", DataType::Utf8, true),
81        Field::new(
82            "_laminar_received_at",
83            DataType::Timestamp(TimeUnit::Nanosecond, None),
84            false,
85        ),
86    ]))
87}
88
89#[cfg(test)]
90mod tests;