Skip to content

Commit ce2fca0

Browse files
committed
rename Event as ParseableEvent
1 parent 5c67134 commit ce2fca0

File tree

1 file changed

+3
-2
lines changed

1 file changed

+3
-2
lines changed

src/connectors/kafka/processor.rs

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ use crate::connectors::kafka::config::BufferConfig;
2121
use crate::connectors::kafka::{ConsumerRecord, StreamConsumer, TopicPartition};
2222
use crate::event::format;
2323
use crate::event::format::EventFormat;
24+
use crate::event::Event as ParseableEvent;
2425
use crate::handlers::http::ingest::create_stream_if_not_exists;
2526
use crate::metadata::{SchemaVersion, STREAM_INFO};
2627
use crate::storage::StreamType;
@@ -41,7 +42,7 @@ impl ParseableSinkProcessor {
4142
async fn deserialize(
4243
&self,
4344
consumer_record: &ConsumerRecord,
44-
) -> anyhow::Result<Option<crate::event::Event>> {
45+
) -> anyhow::Result<Option<ParseableEvent>> {
4546
let stream_name = consumer_record.topic.as_str();
4647

4748
create_stream_if_not_exists(stream_name, &StreamType::UserDefined.to_string()).await?;
@@ -69,7 +70,7 @@ impl ParseableSinkProcessor {
6970
let (record_batch, is_first) =
7071
event.into_recordbatch(&schema, None, None, SchemaVersion::V1)?;
7172

72-
let p_event = crate::event::Event {
73+
let p_event = ParseableEvent {
7374
rb: record_batch,
7475
stream_name: stream_name.to_string(),
7576
origin_format: "json",

0 commit comments

Comments
 (0)