Flink fraud detect
https://github.com/afedulov/fraud-detection-demo
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18
| DataStream<Transaction> transactions = getTransactionsStream(env);
DataStream<Rule> rulesUpdateStream =; BroadcastStream<Rule> rulesStream = rulesUpdateStream.broadcast(Descriptors.rulesDescriptor);
DataStream<Alert> alerts = transactions .connect(rulesStream) .process(new DynamicKeyFunction()) .uid("DynamicKeyFunction") .name("Dynamic Partitioning Function") .keyBy((keyed) -> keyed.getKey()) .connect(rulesStream) .process(new DynamicAlertFunction()) .uid("DynamicAlertFunction") .name("Dynamic Rule Evaluation Function");
|
metric
1 2
| ruleCounterGauge = new RuleCounterGauge(); getRuntimeContext().getMetricGroup().gauge("numberOfActiveRules", ruleCounterGauge);
|