0%

Flink欺诈检测

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);