-
Notifications
You must be signed in to change notification settings - Fork 16
/
Copy pathreceive_super_stream.rs
64 lines (59 loc) · 2.12 KB
/
receive_super_stream.rs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
use futures::StreamExt;
use rabbitmq_stream_client::error::StreamCreateError;
use rabbitmq_stream_client::types::{
ByteCapacity, OffsetSpecification, ResponseCode, SuperStreamConsumer,
};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
use rabbitmq_stream_client::Environment;
let environment = Environment::builder().build().await?;
let message_count = 100_000;
let super_stream = "hello-rust-super-stream";
let create_response = environment
.stream_creator()
.max_length(ByteCapacity::GB(5))
.create_super_stream(super_stream, 3, None)
.await;
if let Err(e) = create_response {
if let StreamCreateError::Create { stream, status } = e {
match status {
// we can ignore this error because the stream already exists
ResponseCode::StreamAlreadyExists => {}
err => {
println!("Error creating stream: {:?} {:?}", stream, err);
}
}
}
}
println!(
"Super stream consumer example, consuming messages from the super stream {}",
super_stream
);
let mut super_stream_consumer: SuperStreamConsumer = environment
.super_stream_consumer()
.offset(OffsetSpecification::First)
.client_provided_name("my super stream consumer for hello rust")
.build(super_stream)
.await
.unwrap();
for _ in 0..message_count {
let delivery = super_stream_consumer.next().await.unwrap();
{
let delivery = delivery.unwrap();
println!(
"Got message: {:#?} from stream: {} with offset: {}",
delivery
.message()
.data()
.map(|data| String::from_utf8(data.to_vec()).unwrap())
.unwrap(),
delivery.stream(),
delivery.offset()
);
}
}
println!("Stopping super stream consumer...");
let _ = super_stream_consumer.handle().close().await;
println!("Super stream consumer stopped");
Ok(())
}