|
|
|
@ -1,5 +1,7 @@ |
|
|
|
package com.qniao.iot.machine.event.generator.job; |
|
|
|
|
|
|
|
import com.qniao.iot.machine.event.MachineIotDataReceivedEvent; |
|
|
|
import com.qniao.iot.machine.event.MachineIotDataReceivedEventKafkaDeserializationSchema; |
|
|
|
import org.apache.flink.api.java.utils.ParameterTool; |
|
|
|
import org.apache.flink.connector.kafka.source.KafkaSource; |
|
|
|
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer; |
|
|
|
@ -17,7 +19,7 @@ public class IotMachineEventGeneratorJob { |
|
|
|
.setTopics("root_cloud_iot_report_data_event") |
|
|
|
.setGroupId("root_cloud_iot_data_etl") |
|
|
|
.setStartingOffsets(OffsetsInitializer.earliest()) |
|
|
|
.setValueOnlyDeserializer(new RootCloudIotDataReceiptedEventDeserializationSchema()) |
|
|
|
.setValueOnlyDeserializer(new MachineIotDataReceivedEventKafkaDeserializationSchema()) |
|
|
|
.build(); |
|
|
|
} |
|
|
|
} |