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

Function main

examples/pulsar_processor.rs:5–67  ·  view source on GitHub ↗
()

Source from the content-addressed store, hash-verified

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

Callers

nothing calls this directly

Calls 1

run_topology_1Function · 0.85

Tested by

no test coverage detected