We read every piece of feedback, and take your input very seriously.
To see all available qualifiers, see our documentation.
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Following dataflow panic when a worker which is not in peers broadcast a batch
conf.set_workers(2); let mut results = pegasus::run(conf, || { |input, output| { let worker_id = input.get_worker_index(); let workers = pegasus::get_current_worker().total_peers() as u64; let stream = input.input_from(vec![worker_id as u64])?; stream .flat_map(|id| { Ok(Some(id).into_iter()) })? .repartition(move |id| Ok(*id % (workers/2))) .filter_map(|source| { Ok(Some(source)) })? .broadcast() .sink_into(output) } }).expect("run job fail;"); while let Some(next) = results.next() { let n = next.unwrap(); println!("{}", n); }
The text was updated successfully, but these errors were encountered:
[pegasus] fix bug in (#1954) (#1958)
7f0f053
* [pegasus] fix bug in (#1954l) * tidy up;
bmmcq
No branches or pull requests
Following dataflow panic when a worker which is not in peers broadcast a batch
The text was updated successfully, but these errors were encountered: