pub struct EventWriter { /* private fields */ }
Expand description

Write events exactly once to a given stream.

EventWriter spawns a Reactor that runs in the background for processing incoming events. The write method sends the event to the Reactor asynchronously and returns a tokio::oneshot::Receiver which contains the result of the write to the caller. The maximum size of the serialized event supported is 8MB, writing size larger than that will return an error.

Backpressure

Write has a backpressure mechanism. Internally, it uses Channel to send event to Reactor for processing. Channel has a limited capacity, when its capacity is reached, any further write will not be accepted until enough space has been freed in the Channel.

Retry

The EventWriter implementation provides retry logic to handle connection failures and service host failures. Internal retries will not violate the exactly once semantic so it is better to rely on them than to wrap this with custom retry logic.

Examples

use pravega_client_config::ClientConfigBuilder;
use pravega_client::client_factory::ClientFactory;
use pravega_client_shared::ScopedStream;

#[tokio::main]
async fn main() {
    // assuming Pravega controller is listening at endpoint `localhost:9090`
    let config = ClientConfigBuilder::default()
        .controller_uri("localhost:9090")
        .build()
        .expect("creating config");

    let client_factory = ClientFactory::new(config);

    // assuming scope:myscope and stream:mystream has been created before.
    let stream = ScopedStream::from("myscope/mystream");

    let mut event_writer = client_factory.create_event_writer(stream);

    let payload = "hello world".to_string().into_bytes();
    let result = event_writer.write_event(payload).await;

    assert!(result.await.is_ok())
}

Implementations§

source§

impl EventWriter

source

pub async fn write_event( &mut self, event: Vec<u8> ) -> Receiver<Result<(), Error>>

Write an event without routing key.

A random routing key will be generated in this case.

This method sends the payload to a background task called reactor to process, so the success of this method only means the payload has been sent to the reactor. Applications may want to call await on the returned tokio oneshot to check whether the payload is successfully persisted or not. If oneshot returns an error indicating something is going wrong on the server side, then subsequent calls are also likely to fail.

Examples
let mut event_writer = client_factory.create_event_writer(stream);
// result is a tokio oneshot
let result = event_writer.write_event(payload).await;
result.await.expect("flush to server");
source

pub async fn write_event_by_routing_key( &mut self, routing_key: String, event: Vec<u8> ) -> Receiver<Result<(), Error>>

Writes an event with a routing key.

Same as the write_event.

source

pub async fn flush(&mut self) -> Result<(), Error>

Flush data.

It will wait until all pending appends have acknowledgment.

Trait Implementations§

source§

impl Drop for EventWriter

source§

fn drop(&mut self)

Executes the destructor for this type. Read more

Auto Trait Implementations§

Blanket Implementations§

source§

impl<T> Any for T
where T: 'static + ?Sized,

source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
source§

impl<T> Borrow<T> for T
where T: ?Sized,

source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
source§

impl<T> From<T> for T

source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> Instrument for T

§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided [Span], returning an Instrumented wrapper. Read more
§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
source§

impl<T> Instrument for T

source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
source§

impl<T, U> Into<U> for T
where U: From<T>,

source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

source§

impl<T> IntoRequest<T> for T

source§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
§

impl<T> Pointable for T

§

const ALIGN: usize = _

The alignment of pointer.
§

type Init = T

The type for initializers.
§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
source§

impl<T> Same for T

§

type Output = T

Should always be Self
source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

§

type Error = Infallible

The type returned in the event of a conversion error.
source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a [WithDispatch] wrapper. Read more
§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a [WithDispatch] wrapper. Read more
source§

impl<T> WithSubscriber for T

source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more