pub struct KafkaEventReader<'w, 's, T: BusEvent + Event> { /* private fields */ }Expand description
Kafka-specific BusEventReader implementation that can optionally perform manual commits
and expose consumer lag metrics.
Implementations§
Source§impl<'w, 's, T: BusEvent + Event> KafkaEventReader<'w, 's, T>
impl<'w, 's, T: BusEvent + Event> KafkaEventReader<'w, 's, T>
Sourcepub fn read<C: EventBusConfig>(&mut self, config: &C) -> Vec<EventWrapper<T>>
pub fn read<C: EventBusConfig>(&mut self, config: &C) -> Vec<EventWrapper<T>>
Read all messages for the supplied Kafka configuration. When manual commit is requested (auto commit disabled) the reader ensures the backend has manual commit mode enabled for the consumer group before returning events.
Sourcepub fn commit(
&mut self,
config: &KafkaConsumerConfig,
event: &EventWrapper<T>,
) -> Result<(), KafkaReaderError>
pub fn commit( &mut self, config: &KafkaConsumerConfig, event: &EventWrapper<T>, ) -> Result<(), KafkaReaderError>
Commit the supplied event using Kafka manual offsets. The configuration used to read the event must have auto commit disabled.
Sourcepub fn consumer_lag(
&self,
config: &KafkaConsumerConfig,
) -> Result<HashMap<String, i64>, KafkaReaderError>
pub fn consumer_lag( &self, config: &KafkaConsumerConfig, ) -> Result<HashMap<String, i64>, KafkaReaderError>
Fetch consumer lag per topic for the supplied configuration.
Trait Implementations§
Source§impl<'w, 's, T: BusEvent + Event> BusEventReader<T> for KafkaEventReader<'w, 's, T>
impl<'w, 's, T: BusEvent + Event> BusEventReader<T> for KafkaEventReader<'w, 's, T>
Source§fn read<C: EventBusConfig>(&mut self, config: &C) -> Vec<EventWrapper<T>>
fn read<C: EventBusConfig>(&mut self, config: &C) -> Vec<EventWrapper<T>>
Drain the buffered events for the supplied configuration.
Source§impl<T: BusEvent + Event> SystemParam for KafkaEventReader<'_, '_, T>
impl<T: BusEvent + Event> SystemParam for KafkaEventReader<'_, '_, T>
Source§type Item<'w, 's> = KafkaEventReader<'w, 's, T>
type Item<'w, 's> = KafkaEventReader<'w, 's, T>
The item type returned when constructing this system param.
The value of this associated type should be
Self, instantiated with new lifetimes. Read moreSource§fn init_state(world: &mut World, system_meta: &mut SystemMeta) -> Self::State
fn init_state(world: &mut World, system_meta: &mut SystemMeta) -> Self::State
Registers any
World access used by this SystemParam
and creates a new instance of this param’s State.Source§unsafe fn new_archetype(
state: &mut Self::State,
archetype: &Archetype,
system_meta: &mut SystemMeta,
)
unsafe fn new_archetype( state: &mut Self::State, archetype: &Archetype, system_meta: &mut SystemMeta, )
For the specified
Archetype, registers the components accessed by this SystemParam (if applicable).a Read moreSource§fn apply(state: &mut Self::State, system_meta: &SystemMeta, world: &mut World)
fn apply(state: &mut Self::State, system_meta: &SystemMeta, world: &mut World)
Applies any deferred mutations stored in this
SystemParam’s state.
This is used to apply Commands during ApplyDeferred.Source§fn queue(
state: &mut Self::State,
system_meta: &SystemMeta,
world: DeferredWorld<'_>,
)
fn queue( state: &mut Self::State, system_meta: &SystemMeta, world: DeferredWorld<'_>, )
Queues any deferred mutations to be applied at the next
ApplyDeferred.Source§unsafe fn validate_param<'w, 's>(
state: &'s Self::State,
_system_meta: &SystemMeta,
_world: UnsafeWorldCell<'w>,
) -> Result<(), SystemParamValidationError>
unsafe fn validate_param<'w, 's>( state: &'s Self::State, _system_meta: &SystemMeta, _world: UnsafeWorldCell<'w>, ) -> Result<(), SystemParamValidationError>
Source§unsafe fn get_param<'w, 's>(
state: &'s mut Self::State,
system_meta: &SystemMeta,
world: UnsafeWorldCell<'w>,
change_tick: Tick,
) -> Self::Item<'w, 's>
unsafe fn get_param<'w, 's>( state: &'s mut Self::State, system_meta: &SystemMeta, world: UnsafeWorldCell<'w>, change_tick: Tick, ) -> Self::Item<'w, 's>
Creates a parameter to be passed into a
SystemParamFunction. Read moreimpl<'w, 's, T: BusEvent + Event> ReadOnlySystemParam for KafkaEventReader<'w, 's, T>where
Local<'s, Vec<EventWrapper<T>>>: ReadOnlySystemParam,
Option<ResMut<'w, DrainedTopicMetadata>>: ReadOnlySystemParam,
Local<'s, HashMap<String, usize>>: ReadOnlySystemParam,
Option<Res<'w, KafkaCommitQueue>>: ReadOnlySystemParam,
Option<Res<'w, KafkaLagCacheResource>>: ReadOnlySystemParam,
Option<Res<'w, ProvisionedTopology>>: ReadOnlySystemParam,
Local<'s, HashSet<String>>: ReadOnlySystemParam,
EventWriter<'w, EventBusError<T>>: ReadOnlySystemParam,
Auto Trait Implementations§
impl<'w, 's, T> Freeze for KafkaEventReader<'w, 's, T>
impl<'w, 's, T> !RefUnwindSafe for KafkaEventReader<'w, 's, T>
impl<'w, 's, T> Send for KafkaEventReader<'w, 's, T>
impl<'w, 's, T> Sync for KafkaEventReader<'w, 's, T>
impl<'w, 's, T> Unpin for KafkaEventReader<'w, 's, T>
impl<'w, 's, T> !UnwindSafe for KafkaEventReader<'w, 's, T>
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more
Source§impl<T> Downcast for Twhere
T: Any,
impl<T> Downcast for Twhere
T: Any,
Source§fn into_any(self: Box<T>) -> Box<dyn Any>
fn into_any(self: Box<T>) -> Box<dyn Any>
Converts
Box<dyn Trait> (where Trait: Downcast) to Box<dyn Any>, which can then be
downcast into Box<dyn ConcreteType> where ConcreteType implements Trait.Source§fn into_any_rc(self: Rc<T>) -> Rc<dyn Any>
fn into_any_rc(self: Rc<T>) -> Rc<dyn Any>
Converts
Rc<Trait> (where Trait: Downcast) to Rc<Any>, which can then be further
downcast into Rc<ConcreteType> where ConcreteType implements Trait.Source§fn as_any(&self) -> &(dyn Any + 'static)
fn as_any(&self) -> &(dyn Any + 'static)
Converts
&Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot
generate &Any’s vtable from &Trait’s.Source§fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)
fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)
Converts
&mut Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot
generate &mut Any’s vtable from &mut Trait’s.Source§impl<T> DowncastSend for T
impl<T> DowncastSend for T
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self>
fn into_either(self, into_left: bool) -> Either<Self, Self>
Converts
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
Converts
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more