delta_kernel/metrics/
metered_json.rs1use 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
16pub struct MeteredJsonHandler {
20 inner: Arc<dyn JsonHandler>,
21}
22
23impl MeteredJsonHandler {
24 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}