Skip to main content

delta_kernel/metrics/
metered_json.rs

1//! [`MeteredJsonHandler`] wraps any [`JsonHandler`] so its `read_json_files` emits the
2//! kernel's standard `JsonReadCompleted` span, carrying `(num_files, bytes_read)` exactly
3//! once when the returned iterator is exhausted or dropped. `parse_json` and
4//! `write_json_file` pass through.
5
6use std::sync::Arc;
7
8use crate::metrics::events::emit_json_read_completed;
9use crate::metrics::PrecountedMetricsIterator;
10use crate::schema::SchemaRef;
11use crate::{
12    DeltaResult, DeltaResultIterator, EngineData, FileDataReadResultIterator, FileMeta,
13    FilteredEngineData, JsonHandler, PredicateRef,
14};
15
16/// Decorator over an engine-provided `Arc<dyn JsonHandler>` that emits a
17/// `JsonReadCompleted` span on every `read_json_files` call. `parse_json` and
18/// `write_json_file` are pass-through and emit nothing.
19pub struct MeteredJsonHandler {
20    inner: Arc<dyn JsonHandler>,
21}
22
23impl MeteredJsonHandler {
24    /// Wrap `inner`. Debug-asserts that `inner` is not already a [`MeteredJsonHandler`]
25    /// so spans are emitted exactly once.
26    pub fn new(inner: Arc<dyn JsonHandler>) -> Self {
27        debug_assert!(
28            !inner.any_ref().is::<MeteredJsonHandler>(),
29            "MeteredJsonHandler wraps another MeteredJsonHandler; \
30             remove the outer wrap to avoid double-counting metrics",
31        );
32        Self { inner }
33    }
34}
35
36impl std::fmt::Debug for MeteredJsonHandler {
37    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
38        f.debug_struct("MeteredJsonHandler").finish_non_exhaustive()
39    }
40}
41
42impl JsonHandler for MeteredJsonHandler {
43    fn parse_json(
44        &self,
45        json_strings: Box<dyn EngineData>,
46        output_schema: SchemaRef,
47    ) -> DeltaResult<Box<dyn EngineData>> {
48        self.inner.parse_json(json_strings, output_schema)
49    }
50
51    fn read_json_files(
52        &self,
53        files: &[FileMeta],
54        physical_schema: SchemaRef,
55        predicate: Option<PredicateRef>,
56    ) -> DeltaResult<FileDataReadResultIterator> {
57        let num_files = files.len() as u64;
58        let bytes_read = files.iter().map(|f| f.size).sum();
59        let inner = self
60            .inner
61            .read_json_files(files, physical_schema, predicate)?;
62        Ok(Box::new(PrecountedMetricsIterator::new(
63            inner,
64            num_files,
65            bytes_read,
66            emit_json_read_completed,
67        )))
68    }
69
70    fn write_json_file(
71        &self,
72        path: &url::Url,
73        data: DeltaResultIterator<'_, FilteredEngineData>,
74        overwrite: bool,
75    ) -> DeltaResult<()> {
76        self.inner.write_json_file(path, data, overwrite)
77    }
78}
79
80#[cfg(test)]
81mod tests {
82    use std::sync::Arc;
83
84    use url::Url;
85
86    use super::*;
87    use crate::arrow::array::RecordBatch;
88    use crate::arrow::datatypes::Schema;
89    use crate::engine::arrow_data::ArrowEngineData;
90    use crate::metrics::MetricEvent;
91    use crate::schema::{DataType, StructField, StructType};
92    use crate::utils::test_utils::{install_thread_local_metrics_reporter, CapturingReporter};
93
94    #[derive(Debug)]
95    struct StubJsonHandler;
96
97    fn empty_batch() -> Box<dyn EngineData> {
98        Box::new(ArrowEngineData::new(RecordBatch::new_empty(Arc::new(
99            Schema::empty(),
100        ))))
101    }
102
103    impl JsonHandler for StubJsonHandler {
104        fn parse_json(
105            &self,
106            _json_strings: Box<dyn EngineData>,
107            _output_schema: SchemaRef,
108        ) -> DeltaResult<Box<dyn EngineData>> {
109            Ok(empty_batch())
110        }
111
112        fn read_json_files(
113            &self,
114            _files: &[FileMeta],
115            _physical_schema: SchemaRef,
116            _predicate: Option<PredicateRef>,
117        ) -> DeltaResult<FileDataReadResultIterator> {
118            Ok(Box::new(std::iter::empty()))
119        }
120
121        fn write_json_file(
122            &self,
123            _path: &Url,
124            _data: DeltaResultIterator<'_, FilteredEngineData>,
125            _overwrite: bool,
126        ) -> DeltaResult<()> {
127            Ok(())
128        }
129    }
130
131    fn fake_file(name: &str, size: u64) -> FileMeta {
132        FileMeta {
133            location: Url::parse(&format!("memory:///_delta_log/{name}")).unwrap(),
134            last_modified: 0,
135            size,
136        }
137    }
138
139    fn install_capture() -> (Arc<CapturingReporter>, tracing::subscriber::DefaultGuard) {
140        let reporter = Arc::new(CapturingReporter::default());
141        let guard = install_thread_local_metrics_reporter(reporter.clone());
142        (reporter, guard)
143    }
144
145    fn delta_schema() -> SchemaRef {
146        Arc::new(StructType::try_new([StructField::nullable("x", DataType::INTEGER)]).unwrap())
147    }
148
149    #[test]
150    fn read_json_files_emits_json_read_completed() {
151        let (reporter, _guard) = install_capture();
152        let inner: Arc<dyn JsonHandler> = Arc::new(StubJsonHandler);
153        let handler = MeteredJsonHandler::new(inner);
154
155        let files = vec![fake_file("0.json", 100), fake_file("1.json", 50)];
156        let iter = handler
157            .read_json_files(&files, delta_schema(), None)
158            .unwrap();
159        let _: Vec<_> = iter.collect();
160
161        let events = reporter.events();
162        let read = events
163            .iter()
164            .find(|e| matches!(e, MetricEvent::JsonReadCompleted(_)))
165            .expect("expected JsonReadCompleted event");
166        let MetricEvent::JsonReadCompleted(e) = read else {
167            unreachable!();
168        };
169        assert_eq!(e.num_files, 2);
170        assert_eq!(e.bytes_read, 150);
171    }
172
173    #[test]
174    fn read_json_files_emits_on_drop_without_consumption() {
175        let (reporter, _guard) = install_capture();
176        let inner: Arc<dyn JsonHandler> = Arc::new(StubJsonHandler);
177        let handler = MeteredJsonHandler::new(inner);
178
179        let files = vec![fake_file("0.json", 100), fake_file("1.json", 50)];
180        {
181            let _iter = handler
182                .read_json_files(&files, delta_schema(), None)
183                .unwrap();
184        }
185
186        let events = reporter.events();
187        let read = events
188            .iter()
189            .find(|e| matches!(e, MetricEvent::JsonReadCompleted(_)))
190            .expect("expected JsonReadCompleted event on drop");
191        let MetricEvent::JsonReadCompleted(e) = read else {
192            unreachable!();
193        };
194        assert_eq!(e.num_files, 2);
195        assert_eq!(e.bytes_read, 150);
196    }
197
198    #[test]
199    fn read_json_files_emits_zero_event_for_empty_input() {
200        let (reporter, _guard) = install_capture();
201        let inner: Arc<dyn JsonHandler> = Arc::new(StubJsonHandler);
202        let handler = MeteredJsonHandler::new(inner);
203
204        let iter = handler.read_json_files(&[], delta_schema(), None).unwrap();
205        let _: Vec<_> = iter.collect();
206
207        let events = reporter.events();
208        let read = events
209            .iter()
210            .find(|e| matches!(e, MetricEvent::JsonReadCompleted(_)))
211            .expect("expected zero-valued JsonReadCompleted");
212        let MetricEvent::JsonReadCompleted(e) = read else {
213            unreachable!();
214        };
215        assert_eq!(e.num_files, 0);
216        assert_eq!(e.bytes_read, 0);
217    }
218
219    #[test]
220    #[should_panic(expected = "wraps another MeteredJsonHandler")]
221    fn new_panics_on_double_wrap() {
222        let inner: Arc<dyn JsonHandler> = Arc::new(StubJsonHandler);
223        let once: Arc<dyn JsonHandler> = Arc::new(MeteredJsonHandler::new(inner));
224        let _twice = MeteredJsonHandler::new(once);
225    }
226}