David Wisecup
06/12/2023, 7:42 PMDavid Anderson
06/12/2023, 7:45 PMDavid Wisecup
06/12/2023, 8:14 PMDavid Wisecup
06/12/2023, 8:31 PMWatermarkStrategy.<Transaction>forBoundedOutOfOrderness(Duration.ofSeconds(10))David Wisecup
06/13/2023, 1:32 AMDataStreamSource<Transaction> kafkaStream =
env.fromSource(source, WatermarkStrategy
.<Transaction>forBoundedOutOfOrderness(Duration.ofSeconds(1))
.withTimestampAssigner(new SerializableTimestampAssigner<Transaction>() {
@Override
public long extractTimestamp(Transaction txn, long recordTimestamp) {
if (txn != null && txn.getDate() != null) {
Date date = txn.getDate();
// System.out.println("Extracting timestamp: " + timestamp);
return date.toInstant().toEpochMilli();
}
// fallback to stream time
String pCode = txn.get(Transaction.PCODE);
System.out.printf("Unable to extract date from txn: %s. Defaulting to recordTimestamp\n", pCode);
return recordTimestamp;
}
}), "Kafka Source With Timestamps and Watermarks");David Anderson
06/13/2023, 9:24 AM.withIdleness or by reducing the parallelism.David Wisecup
06/13/2023, 3:16 PMDavid Wisecup
06/13/2023, 3:21 PMDavid Anderson
06/13/2023, 8:20 PM