This commit is contained in:
25
src/main/kotlin/mq/ConsumerWrapper.kt
Normal file
25
src/main/kotlin/mq/ConsumerWrapper.kt
Normal file
@ -0,0 +1,25 @@
|
||||
package mq
|
||||
|
||||
import com.rabbitmq.client.BuiltinExchangeType
|
||||
import com.rabbitmq.client.ConnectionFactory
|
||||
import config.EnvConfig
|
||||
|
||||
class ConsumerWrapper {
|
||||
private val envConfig = EnvConfig()
|
||||
|
||||
fun recieve(){
|
||||
val factory = ConnectionFactory()
|
||||
factory.host = envConfig.mqHost
|
||||
factory.username = envConfig.mqUserName
|
||||
factory.password = envConfig.mqPassWord
|
||||
|
||||
val inputConnection = factory.newConnection()
|
||||
val inputChannel = inputConnection.createChannel()
|
||||
|
||||
inputChannel.exchangeDeclare(envConfig.mqExchange, BuiltinExchangeType.FANOUT)
|
||||
val inputQueueName = inputChannel.queueDeclare().queue
|
||||
inputChannel.queueBind(inputQueueName, envConfig.mqExchange, "")
|
||||
|
||||
inputChannel.basicConsume(inputQueueName, true, DatabaseConsumer())
|
||||
}
|
||||
}
|
37
src/main/kotlin/mq/DatabaseConsumer.kt
Normal file
37
src/main/kotlin/mq/DatabaseConsumer.kt
Normal file
@ -0,0 +1,37 @@
|
||||
package mq
|
||||
|
||||
import api.ApiObject
|
||||
import com.google.gson.Gson
|
||||
import com.rabbitmq.client.AMQP
|
||||
import com.rabbitmq.client.Consumer
|
||||
import com.rabbitmq.client.Envelope
|
||||
import com.rabbitmq.client.ShutdownSignalException
|
||||
import database.service.IResultObjectService
|
||||
import org.koin.core.component.KoinComponent
|
||||
import org.koin.core.component.inject
|
||||
|
||||
class DatabaseConsumer: Consumer, KoinComponent {
|
||||
private val resultObjectService : IResultObjectService by inject()
|
||||
private val gson = Gson()
|
||||
override fun handleConsumeOk(consumerTag : String?) {
|
||||
}
|
||||
override fun handleCancelOk(p0 : String?) {
|
||||
throw UnsupportedOperationException()
|
||||
}
|
||||
override fun handleRecoverOk(p0 : String?) {
|
||||
throw UnsupportedOperationException()
|
||||
}
|
||||
override fun handleCancel(p0 : String?) {
|
||||
throw UnsupportedOperationException()
|
||||
}
|
||||
|
||||
override fun handleShutdownSignal(consumerTag: String?, sig: ShutdownSignalException?) {
|
||||
println("got shutdown signal")
|
||||
}
|
||||
|
||||
override fun handleDelivery(consumerTag : String?, envelope : Envelope?, basicProperties : AMQP.BasicProperties?, body : ByteArray?) {
|
||||
val rawJson = body!!.toString(Charsets.UTF_8)
|
||||
val apiObject = gson.fromJson(rawJson, ApiObject::class.java)
|
||||
resultObjectService.addOne(apiObject)
|
||||
}
|
||||
}
|
Reference in New Issue
Block a user