Lib update test
This commit is contained in:
parent
894d8717c0
commit
f245f14291
@ -23,7 +23,7 @@ repositories {
|
|||||||
}
|
}
|
||||||
|
|
||||||
dependencies {
|
dependencies {
|
||||||
implementation("no.iktdev.streamit.library:streamit-library-kafka:0.0.2-alpha23")
|
implementation("no.iktdev.streamit.library:streamit-library-kafka:0.0.2-alpha26")
|
||||||
implementation("no.iktdev:exfl:0.0.4-SNAPSHOT")
|
implementation("no.iktdev:exfl:0.0.4-SNAPSHOT")
|
||||||
|
|
||||||
implementation("com.github.pgreze:kotlin-process:1.3.1")
|
implementation("com.github.pgreze:kotlin-process:1.3.1")
|
||||||
|
|||||||
@ -3,3 +3,4 @@ rootProject.name = "Reader"
|
|||||||
include(":CommonCode")
|
include(":CommonCode")
|
||||||
project(":CommonCode").projectDir = File("../CommonCode")
|
project(":CommonCode").projectDir = File("../CommonCode")
|
||||||
|
|
||||||
|
include(":streamit-library-kafka")
|
||||||
@ -29,17 +29,7 @@ class StreamsReader {
|
|||||||
|
|
||||||
|
|
||||||
init {
|
init {
|
||||||
val listener = StreamReaderListener(CommonConfig.kafkaTopic, defaultConsumer, listOf(EVENT_READER_RECEIVED_FILE.event))
|
object: SimpleMessageListener(CommonConfig.kafkaTopic, defaultConsumer, listOf(EVENT_READER_RECEIVED_FILE.event)) {
|
||||||
listener.listen()
|
|
||||||
|
|
||||||
/*object: EventMessageListener(CommonConfig.kafkaTopic, defaultConsumer, listOf(EVENT_READER_RECEIVED_FILE.event)) {
|
|
||||||
override fun onMessage(data: ConsumerRecord<String, Message>) {
|
|
||||||
|
|
||||||
}
|
|
||||||
}.listen()*/
|
|
||||||
}
|
|
||||||
|
|
||||||
inner class StreamReaderListener(topic: String, consumer: DefaultConsumer, accepts: List<String>): SimpleMessageListener(topic = topic, consumer = consumer, accepts = accepts) {
|
|
||||||
override fun onMessage(data: ConsumerRecord<String, Message>) {
|
override fun onMessage(data: ConsumerRecord<String, Message>) {
|
||||||
logger.info { "RECORD: ${data.key()}" }
|
logger.info { "RECORD: ${data.key()}" }
|
||||||
logger.info { "Active filters: ${this.accepts.joinToString(",") }}" }
|
logger.info { "Active filters: ${this.accepts.joinToString(",") }}" }
|
||||||
@ -81,6 +71,7 @@ class StreamsReader {
|
|||||||
val message = Message(status = Status( statusType = if (resultCode == 0) StatusType.SUCCESS else StatusType.ERROR), data = output.joinToString("\n"))
|
val message = Message(status = Status( statusType = if (resultCode == 0) StatusType.SUCCESS else StatusType.ERROR), data = output.joinToString("\n"))
|
||||||
messageProducer.sendMessage(KnownEvents.EVENT_READER_RECEIVED_STREAMS.event, message)
|
messageProducer.sendMessage(KnownEvents.EVENT_READER_RECEIVED_STREAMS.event, message)
|
||||||
}
|
}
|
||||||
|
}.listen()
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
}
|
|
||||||
Loading…
Reference in New Issue
Block a user