MCPcopy Create free account
hub / github.com/MaterializeInc/materialize / run_worker

Method run_worker

src/storage/src/server.rs:123–204  ·  view source on GitHub ↗
(
        &self,
        timely_worker: &mut TimelyWorker,
        client_rx: mpsc::UnboundedReceiver<(
            Uuid,
            mpsc::UnboundedReceiver<StorageCommand>,
            mpsc::Unbound

Source from the content-addressed store, hash-verified

121 const NAME: &str = "storage";
122
123 fn run_worker(
124 &self,
125 timely_worker: &mut TimelyWorker,
126 client_rx: mpsc::UnboundedReceiver<(
127 Uuid,
128 mpsc::UnboundedReceiver<StorageCommand>,
129 mpsc::UnboundedSender<StorageResponse>,
130 )>,
131 ) {
132 // Register a timely logger that forwards events to the compute logging dataflow.
133 // Assign by local worker index so storage worker x matches compute worker x.
134 let local_index = timely_worker.index() % self.workers_per_process;
135 let writer = self.timely_log_writers.lock().unwrap()[local_index].take();
136 if let Some(writer) = writer {
137 use timely::dataflow::operators::capture::{Event, EventPusher};
138 use timely::logging::TimelyEventBuilder;
139
140 // We use an approach similar to compute's logging: wrap the writer in
141 // a BatchLogger that translates Logger callbacks into Event pushes,
142 // then register the Logger with timely's log_register.
143 let interval_ms = 1000u128; // 1 second batching interval
144 let mut time_ms = mz_repr::Timestamp::from(0u64);
145 let mut event_pusher = writer;
146 let now = std::time::Instant::now();
147 let start_offset = std::time::SystemTime::now()
148 .duration_since(std::time::SystemTime::UNIX_EPOCH)
149 .expect("Failed to get duration since Unix epoch");
150
151 let logger = timely::logging_core::Logger::<TimelyEventBuilder>::new(
152 now,
153 start_offset,
154 move |time: &std::time::Duration,
155 data: &mut Option<Vec<(std::time::Duration, TimelyEvent)>>| {
156 if let Some(mut data) = data.take() {
157 // Filter park/unpark events and remap IDs before handing events
158 // off to compute. Compute's park tracking assumes a single
159 // timely runtime; mixing in storage's park events would break
160 // it. Remapping ensures storage operator/channel IDs don't
161 // collide with compute's.
162 data.retain_mut(|(_, event)| {
163 if matches!(event, TimelyEvent::Park(_)) {
164 return false;
165 }
166 remap_timely_event_ids(event);
167 true
168 });
169 event_pusher.push(Event::Messages(time_ms, data));
170 } else {
171 // Advance progress.
172 let new_time_ms: u64 = (((time.as_millis() / interval_ms) + 1)
173 * interval_ms)
174 .try_into()
175 .expect("must fit");
176 let new_time_ms = mz_repr::Timestamp::from(new_time_ms);
177 if time_ms < new_time_ms {
178 event_pusher
179 .push(Event::Progress(vec![(new_time_ms, 1), (time_ms, -1)]));
180 time_ms = new_time_ms;

Callers

nothing calls this directly

Calls 13

nowFunction · 0.85
remap_timely_event_idsFunction · 0.85
ProgressFunction · 0.85
cloneFunction · 0.85
unwrapMethod · 0.80
lockMethod · 0.80
expectMethod · 0.80
indexMethod · 0.45
takeMethod · 0.45
pushMethod · 0.45
try_intoMethod · 0.45
runMethod · 0.45

Tested by

no test coverage detected