-
Notifications
You must be signed in to change notification settings - Fork 18
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
* adding producer examples * Update documentation Signed-off-by: Gabriele Santomaggio <[email protected]> --------- Signed-off-by: Gabriele Santomaggio <[email protected]> Co-authored-by: Gabriele Santomaggio <[email protected]>
- Loading branch information
1 parent
30df99f
commit c05e420
Showing
5 changed files
with
269 additions
and
76 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,69 @@ | ||
use futures::StreamExt; | ||
use rabbitmq_stream_client::{ | ||
types::{ByteCapacity, Message, OffsetSpecification}, | ||
Environment, | ||
}; | ||
use tracing::{info, Level}; | ||
use tracing_subscriber::FmtSubscriber; | ||
|
||
#[tokio::main] | ||
async fn main() -> Result<(), Box<dyn std::error::Error>> { | ||
let subscriber = FmtSubscriber::builder() | ||
.with_max_level(Level::TRACE) | ||
.finish(); | ||
|
||
tracing::subscriber::set_global_default(subscriber).expect("setting default subscriber failed"); | ||
let environment = Environment::builder() | ||
.host("localhost") | ||
.port(5552) | ||
.build() | ||
.await?; | ||
|
||
let _ = environment.delete_stream("batch_send").await; | ||
let message_count = 10; | ||
environment | ||
.stream_creator() | ||
.max_length(ByteCapacity::GB(2)) | ||
.create("batch_send") | ||
.await?; | ||
|
||
let producer = environment.producer().build("batch_send").await?; | ||
|
||
let mut messages = Vec::with_capacity(message_count); | ||
for i in 0..message_count { | ||
let msg = Message::builder().body(format!("message{}", i)).build(); | ||
messages.push(msg); | ||
} | ||
|
||
producer | ||
.batch_send(messages, |confirmation_status| async move { | ||
info!("Message confirmed with status {:?}", confirmation_status); | ||
}) | ||
.await?; | ||
|
||
producer.close().await?; | ||
|
||
let mut consumer = environment | ||
.consumer() | ||
.offset(OffsetSpecification::First) | ||
.build("batch_send") | ||
.await | ||
.unwrap(); | ||
|
||
for _ in 0..message_count { | ||
let delivery = consumer.next().await.unwrap()?; | ||
info!( | ||
"Got message : {:?} with offset {}", | ||
delivery | ||
.message() | ||
.data() | ||
.map(|data| String::from_utf8(data.to_vec())), | ||
delivery.offset() | ||
); | ||
} | ||
|
||
consumer.handle().close().await.unwrap(); | ||
|
||
environment.delete_stream("batch_send").await?; | ||
Ok(()) | ||
} |
Oops, something went wrong.