MCPcopy Create free account
hub / github.com/Rustixir/tokio_sky / main

Function main

examples/pulsar_processor_complex.rs:6–77  ·  view source on GitHub ↗
()

Source from the content-addressed store, hash-verified

4
5#[tokio::main]
6async fn main() {
7
8 let producer_factory = || Prod;
9 let producer_concurrency = 3;
10 let producer_router = RouterType::RoundRobin;
11 let producer_buffer_pool = 100;
12
13
14 // Pulsar Config
15 let addr = "pulsar://127.0.0.1:6650";
16 let topic = "non-persistent://public/default/topic1";
17 let pulsar_instance_name = "pulsar_processor";
18
19 let pulsar: Pulsar<_> = Pulsar::builder(addr, TokioExecutor).build().await.unwrap();
20 let opts = producer::ProducerOptions {
21 schema: Some(proto::Schema {
22 r#type: proto::schema::Type::String as i32,
23 ..Default::default()
24 }),
25 ..Default::default()
26 };
27
28
29 let pulsar_proc_factory = || || PulsarProcessor::new(pulsar, opts, topic, pulsar_instance_name).unwrap();
30 let pulsar_proc_concurrency = 1;
31 let pulsar_proc_router = RouterType::RoundRobin;
32 let pulsar_proc_buffer_size = 100;
33
34
35 let result_handler_proc_factory = || DeliveryHandler;
36 let result_handler_proc_concurrency = 1;
37 let result_handler_proc_buffer_size = 10;
38
39
40
41 // 1. create X processor instances by 'proc_concurrency'
42 //
43 // 2. create X producer instances by 'producer_concurrency'
44 //
45 // 3. create topology and syncing
46 //
47
48
49 // / \
50 // producer-1 / \
51 // producer-2 ----- > Pulsar_processor -> processor
52 // producer-3 \ /
53 // \ /
54 //
55
56 let safe_shutdown =
57 run_topology_2(
58 producer_factory,
59 producer_concurrency,
60 producer_router,
61 producer_buffer_pool,
62
63 pulsar_proc_factory,

Callers

nothing calls this directly

Calls 1

run_topology_2Function · 0.85

Tested by

no test coverage detected