Skip to content
Snippets Groups Projects

Migrate Beam benchmark implementation

Compare and Show latest version
4 files
+ 76
22
Compare changes
  • Side-by-side
  • Inline
Files
4
@@ -14,20 +14,19 @@ import titan.ccp.model.records.ActivePowerRecord;
@@ -14,20 +14,19 @@ import titan.ccp.model.records.ActivePowerRecord;
/**
/**
* Simple {@link PTransform} that read from Kafka using {@link KafkaIO}.
* Simple {@link PTransform} that read from Kafka using {@link KafkaIO}.
*/
*/
public class KafkaAggregatedPowerRecordReader extends
public class KafkaActivePowerRecordReader extends
PTransform<PBegin, PCollection<KV<String, ActivePowerRecord>>> {
PTransform<PBegin, PCollection<KV<String, ActivePowerRecord>>> {
private static final long serialVersionUID = 2603286150183186115L;
private static final long serialVersionUID = 2603286150183186115L;
private final PTransform<PBegin, PCollection<KV<String, ActivePowerRecord>>> reader;
private final PTransform<PBegin, PCollection<KV<String, ActivePowerRecord>>> reader;
/**
/**
* Instantiates a {@link PTransform} that reads from Kafka with the given Configuration.
* Instantiates a {@link PTransform} that reads from Kafka with the given Configuration.
*/
*/
@SuppressWarnings({"unchecked", "rawtypes"})
@SuppressWarnings({"unchecked", "rawtypes"})
public KafkaAggregatedPowerRecordReader(final String bootstrapServer, final String inputTopic,
public KafkaActivePowerRecordReader(final String bootstrapServer, final String inputTopic,
final Properties consumerConfig) {
final Properties consumerConfig) {
super();
super();
// Check if boostrap server and inputTopic are defined
// Check if boostrap server and inputTopic are defined
Loading