Appearance
Flink 实时任务怎么关联 IP 归属地维表?异步 IO + 本地缓存给日志打地域标签(Java 实现)
实时风控和实时大盘绕不开一件事:给每条事件打上地域标签。某省登录失败率突然翻倍、某运营商的支付成功率掉底、某个城市的注册量异常——这些告警的前提是流里的每条记录都带着"省 / 市 / 运营商"。而这些信息不在你的业务数据里,只在用户 IP 上。
于是就有了"IP 归属地维表"这个需求:实时给它 join 一份 IP → 地理位置的数据。看着简单,真做起来有三条路,每条都有坑。
一、三条路的坑
第一条:把 IP 库预加载进算子内存。 全量 IP 库是几百万条 CIDR 记录,塞进 TaskManager 内存不现实;就算塞进去了,库每月更新,你得靠 BroadcastState 把新版本广播到所有并发实例,还要处理"切换瞬间两条记录打了不同版本的标签"。工程量大,收益一般。
第二条:Flink SQL 的 lookup join 直连一张 MySQL 维表。 写起来最舒服,但每条流记录都会发一次维表查询。日志里 IP 的重复率极高(同一个 NAT 出口 IP 一秒出现几百次),一条记录一次查询等于把维表打爆。你要往前面套一个在线接口,接口免费额度 60 次/分钟,几秒钟就被打光。
第三条:AsyncDataStream 做异步维表关联,自己管缓存和限速。 这是本文的做法。核心思路:缓存挡掉重复 IP,只对没见过的 IP 发请求;请求在异步算子里飞,不阻塞流;查不到不失败,打「未知」放行。
这篇文章用 Flink 的 AsyncDataStream 实现第三种,代码是完整可跑的,并且我把三次真实翻车都写进去了——包括两核机器上 10 条数据全部被打成"未知"的那次。
二、方案设计
数据流长这样:
Kafka source → 解析成 AccessEvent(ip, userId, action)
→ AsyncDataStream.unorderedWait(GeoAsyncFunction) ← 打标签
→ keyBy(省) → 累计计数/阈值告警 → Sink几个设计点值得先说清:
并发上限不是越大越好。 unorderedWait 的 capacity 参数表示同时在飞的元素数。设 50 意味着最多 50 条记录在等结果,超出的会形成背压——这是对的,因为下游接口有额度限制,压不住就得让上游慢下来。
超时必须兜底。 Flink 的元素超时默认会把作业搞挂(后面有实测报错),而限速场景下超时是常态(一次缓存穿透要等 1.1 秒),所以 AsyncFunction.timeout() 一定要覆写成"打未知标签放行"。
不要在 asyncInvoke 里做阻塞式 HTTP 调用还得靠公共线程池。 这是最容易踩的坑:Flink 的 AsyncFunction 只负责调度,真正的阻塞调用得自己安排线程。用 CompletableFuture.supplyAsync() 不传线程池时,用的是 ForkJoinPool.commonPool(),两核机器上它只有 1 个线程——10 条事件排在一个线程上,每条等 1.1 秒限速,全都在超时线外面,最后 10 条数据全变成"未知"。必须给它一个自己的线程池。
并行度和限速额度是要算的账。 免费版 60 次/分钟是按调用方 IP 算的。作业并行度 2,两个 subtask 各自限速 1.1 秒,出口 IP 只有一个,实际速率就是 2 倍——很容易触发限速(返回 429)。三种解法:把并行度压到 1(吞吐靠缓存)、把令牌桶做成全局的(放 Redis)、或者升 VIP 版(按次计费 0.05 元/千次起,18 万次/分钟)。
三、Java 实现
环境:JDK 17+、Flink 1.20、只用 JDK 自带的 java.net.http.HttpClient 发请求,不引额外 HTTP 和 JSON 依赖。
3.1 事件模型(POJO,略)
java
public static class AccessEvent { // 进来:只有 IP 和业务字段
public final String ip, userId, action;
public AccessEvent(String ip, String userId, String action) { ... }
}
public static class TaggedEvent { // 出去:多了省/市/大区/运营商
public final String ip, userId, action, prov, city, isp, bigArea;
public final boolean unknown; // 查不到的标记,下游规则可以据此区别对待
public String key() { return (prov == null || prov.isEmpty()) ? "未知" : prov; }
}3.2 查询器:LRU 缓存 + 串行限速 + 一次重试
java
public static class GeoLookup {
private static final String API = "https://ip9.com.cn/get";
private static final String UA = "Mozilla/5.0 (X11; Linux x86_64) flink-geo/1.0";
private static final long MIN_INTERVAL_MS = 1100L; // 免费版 60 次/分钟
private final HttpClient client;
private final Map<String, String[]> cache;
private final AtomicLong lastCall = new AtomicLong(0L);
public GeoLookup(int cacheSize) {
this.client = HttpClient.newBuilder()
.connectTimeout(Duration.ofSeconds(3))
.version(HttpClient.Version.HTTP_1_1)
.build();
// accessOrder=true + 覆写 removeEldestEntry 就是一个够用的 LRU
this.cache = Collections.synchronizedMap(new LinkedHashMap<>(1024, 0.75f, true) {
@Override
protected boolean removeEldestEntry(Map.Entry<String, String[]> eldest) {
return size() > cacheSize;
}
});
}
/** 限速必须串行:同一时刻只放一个请求过 1.1 秒的间隔 */
private synchronized void acquire() throws InterruptedException {
long wait = MIN_INTERVAL_MS - (System.currentTimeMillis() - lastCall.get());
if (wait > 0) Thread.sleep(wait);
lastCall.set(System.currentTimeMillis());
}
public String[] lookup(String ip) throws Exception {
String[] hit = cache.get(ip);
if (hit != null) return hit; // 命中缓存,零网络开销
String url = API + "?ip=" + URLEncoder.encode(ip, StandardCharsets.UTF_8);
Exception lastError = null;
for (int attempt = 0; attempt < 2; attempt++) {
acquire();
HttpRequest req = HttpRequest.newBuilder(URI.create(url))
.timeout(Duration.ofSeconds(5))
.header("User-Agent", UA) // 不带 UA 可能被前置防护挡掉
.GET().build();
try {
HttpResponse<String> resp = client.send(req, HttpResponse.BodyHandlers.ofString());
if (resp.statusCode() == 429) { // 触发限速,退避后重试
Thread.sleep(2000L * (attempt + 1));
continue;
}
if (resp.statusCode() != 200) throw new IllegalStateException("HTTP " + resp.statusCode());
String body = resp.body();
if (!"200".equals(field(body, "ret"))) { // 非法 IP 返回 400,此时 data 是空数组
throw new IllegalStateException("ret=" + field(body, "ret"));
}
String[] out = new String[]{field(body, "prov"), field(body, "city"),
field(body, "isp"), field(body, "big_area")};
cache.put(ip, out);
return out;
} catch (Exception e) {
lastError = e;
Thread.sleep(500L * (attempt + 1));
}
}
throw lastError == null ? new IllegalStateException("unknown") : lastError;
}
}field() 是一个十来行的字段提取函数(返回的是扁平 JSON,没必要为它引一个 JSON 库)。有个细节值得记住:运营商名里带 & 的值会返回成 \u0026,不还原转义的话,你打出来的标签就是 AT\u0026T。
3.3 异步算子:缓存命中即返回,超时兜底
java
public static class GeoAsyncFunction extends RichAsyncFunction<AccessEvent, TaggedEvent> {
private final int cacheSize, poolSize;
private transient GeoLookup lookup;
private transient ExecutorService pool;
public GeoAsyncFunction(int cacheSize, int poolSize) {
this.cacheSize = cacheSize; this.poolSize = poolSize;
}
@Override
public void open(Configuration parameters) {
lookup = new GeoLookup(cacheSize);
// 关键:阻塞式 HTTP 调用要放在自己的线程池里
pool = Executors.newFixedThreadPool(poolSize, r -> {
Thread t = new Thread(r, "geo-lookup");
t.setDaemon(true);
return t;
});
}
@Override
public void asyncInvoke(AccessEvent input, ResultFuture<TaggedEvent> resultFuture) {
CompletableFuture
.supplyAsync(() -> {
try {
String[] g = lookup.lookup(input.ip);
return new TaggedEvent(input, g[0], g[1], g[2], g[3], false);
} catch (Exception e) {
// 单条查不到不阻塞流:打「未知」放行,交给下游规则决定
return new TaggedEvent(input, "未知", "", "", "", true);
}
}, pool)
.thenAccept(t -> resultFuture.complete(Collections.singletonList(t)));
}
/** 超时兜底:不覆写它,超时会让整个作业失败重启 */
@Override
public void timeout(AccessEvent input, ResultFuture<TaggedEvent> resultFuture) {
resultFuture.complete(Collections.singletonList(
new TaggedEvent(input, "未知", "", "", "", true)));
}
@Override
public void close() { if (pool != null) pool.shutdownNow(); }
}3.4 拼起来
java
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(2);
DataStream<AccessEvent> events = env.fromElements(
new AccessEvent("114.114.114.114", "u1001", "login"),
new AccessEvent("101.226.4.6", "u1002", "login"),
new AccessEvent("114.114.114.114", "u1003", "pay"),
new AccessEvent("183.204.111.150", "u1004", "login"),
new AccessEvent("8.8.8.8", "u1005", "login"),
new AccessEvent("114.114.114.114", "u1006", "logout"),
new AccessEvent("203.198.0.1", "u1007", "login"),
new AccessEvent("240e:ff:e02c:1:0:ff:b0e4:20f", "u1008", "login"),
new AccessEvent("10.0.0.5", "u1009", "health"),
new AccessEvent("45.153.160.2", "u1010", "login"));
DataStream<TaggedEvent> tagged = AsyncDataStream.unorderedWait(
events,
new GeoAsyncFunction(10_000, 8),
10_000, TimeUnit.MILLISECONDS, // 单条超时:限速场景下别设太小
50); // 并发上限 = 同时在飞的元素数
tagged.print("tagged");
tagged.keyBy(TaggedEvent::key)
.process(new ProvinceCounter()) // 按省累计,state 存 ValueState,重启不丢
.print("counter");真上生产时把 print 换成 Kafka / ClickHouse / Doris 的 sink。注意不要在 asyncInvoke 里直接写外部存储——那等于把异步算子又变回同步阻塞。
四、实跑记录:三次翻车和最终结果
我在一台 2 核的机器上跑了全套,环境是 JDK 21 + Flink 1.20.0 官方二进制包。下面都是真实输出。
翻车一:JDK 21 起不来。 第一次运行直接抛:
text
java.io.IOException: Serializing the source elements failed:
java.lang.reflect.InaccessibleObjectException: Unable to make field private final
java.lang.Object[] java.util.Arrays$ArrayList.a accessible: module java.base
does not "opens java.util" to unnamed module @31206beb
at com.twitter.chill.java.ArraysAsListSerializer.<init>(ArraysAsListSerializer.java:69)Flink 1.20 官方支持到 Java 17,用 21 得自己补模块开放参数:
bash
java --add-opens=java.base/java.util=ALL-UNNAMED \
--add-opens=java.base/java.lang=ALL-UNNAMED \
--add-opens=java.base/java.util.concurrent=ALL-UNNAMED \
-cp "flink/lib/*:out" GeoTagJob翻车二:把超时设成 1 毫秒,作业直接挂。 我把 unorderedWait 的超时临时改成 1ms,想看看超时的行为,结果是:
text
org.apache.flink.runtime.client.JobExecutionException: Job execution failed.
Caused by: java.lang.Exception: Could not complete the stream element:
Record @ (undef) : GeoTagJob$AccessEvent@4324c90f.
at org.apache.flink.streaming.api.operators.async.AsyncWaitOperator$ResultHandler.completeExceptionally
at org.apache.flink.streaming.api.functions.async.AsyncFunction.timeout(AsyncFunction.java:97)也就是说:不覆写 timeout(),元素超时 = 作业失败重启。覆写之后同样的 1ms 超时下,作业不挂了,10 条事件全部带上"未知"标签正常流出。这个兜底在限速场景里是必需的。
翻车三:公共线程池把整批数据打成"未知"。 第一版我用的是 CompletableFuture.supplyAsync(...) 不带线程池,2 核机器上 ForkJoinPool.commonPool() 只有 1 个工作线程,加上我那个串行限速,10 条事件的执行序列是"排队等 1.1 秒 × 9",全在超时线外面:
text
tagged:2> 101.226.4.6 未知
tagged:1> 114.114.114.114 未知
tagged:2> 183.204.111.150 未知
...(10 条全部是"未知")换成自己的 Executors.newFixedThreadPool(8)(守护线程)之后,同一份代码、同一台机器,结果就正常了:
text
tagged:1> 114.114.114.114 江苏 南京 华东 114DNS
tagged:2> 101.226.4.6 上海 上海 华东 DNSpai/电信
tagged:2> 183.204.111.150 河南 新乡 华中 中国移动
tagged:2> 240e:ff:e02c:1:0:ff:b0e4:20f 广东 广州 华南 中国电信
tagged:1> 203.198.0.1 香港 香港 港澳台 香港电讯
tagged:1> 8.8.8.8 Google Cloud
tagged:1> 10.0.0.5 内网地址
tagged:2> 45.153.160.2 AT&T
counter:2> 省份=江苏 累计=1 最新=114.114.114.114(114DNS)
counter:2> 省份=江苏 累计=2 最新=114.114.114.114(114DNS)
counter:2> 省份=江苏 累计=3 最新=114.114.114.114(114DNS)
counter:1> 省份=未知 累计=3 最新=8.8.8.8(Google Cloud)
counter:1> 省份=香港 累计=1 最新=203.198.0.1(香港电讯)几点观察:
- 10 条事件、9 个不同的 IP(
114.114.114.114出现 3 次),实际只发了 9 次请求,2 次是缓存命中。真实日志里这个比例会高得多。 unorderedWait是真的不保序:tagged:1和tagged:2交错,江苏的三条连续输出但未知的三条被分散。要严格保序就把unorderedWait换成orderedWait,代价是慢请求会挡住后面的元素。- 境外和保留地址会掉进"未知":
8.8.8.8是"美国"但prov为空,45.153.160.2是"捷克"同样没有省,10.0.0.5干脆是内网(country=保留、isp=内网地址)。所以key()里"有省取省、没省取国家、内网单独一类"这个回退逻辑必须写全,否则你的实时大盘上会只有一个巨大的"未知"。 - 整个作业从启动到结束 23.4 秒,其中 mini cluster 的启动占了一半多,剩下的基本都是 9 次请求的 1.1 秒限速等待。在限速面前,异步 IO 省下的是并发等待时间,不是额度——额度只能靠缓存和 VIP 版本解决。
五、落地时的几个细节
缓存要不要加 TTL? IP 归属地变化很慢,进程内 LRU(示例里 1 万条)通常够用,而缓存穿透的代价只有一次 1.1 秒。真要精确控制,就给缓存值加时间戳,超过 24 小时重新查——比定时清空缓存温和。
IP 归一化能再省一大截。 同一个 C 段里的地址往往属于同一个城市。把 IPv4 按 /24 归一(1.2.3.0/24)、IPv6 按 /64 归一,缓存命中率还能往上抬一截,代价是相邻段跨城的少数误判。风控、大盘这类场景完全可以接受。
429 和 400 要分开处理。 429 是限速,退避重试有意义;400 是非法 IP,重试一百次也没用,直接当"未知"。把这些错误码抄进监控,限速发生率是你能提前发现"该升 VIP 了"的信号。
别用归属地做精确定位判定。 接口返回的经纬度是城市中心点,不是用户位置;境外 IP 更是只有国家一级。要算距离、画电子围栏,先确认业务能接受城市级误差——这一点在风控规则里尤其重要,误判一个正常用户的成本比漏过一个可疑请求高得多。
要机房识别就用 VIP 字段。 上面 45.153.160.2 打成"捷克"看着像境外用户,其实是云主机出口。免费版只能从 isp 里看出点端倪(写着 Google Cloud、Cloudflare 的基本是机房),要按"ISP 家庭 / 企业 / IDC 机房"分类拦截,得用 VIP 版的 ip_type 和 ip_asn。
代码里查归属地用的接口是 IP9 的免费接口 https://ip9.com.cn/get?ip=<IP>,不传 ip 参数返回调用方自己的归属地,IPv4/IPv6 都支持,免费版 60 次/分钟、无需注册;返回的 prov、city、big_area、isp 拿来打标签刚好够用,字段说明和 VIP/私有化部署的额度都在官网 https://www.ip9.com.cn 上。如果流处理之外还有日志侧的统计需求,可以看本站《用IP查询分析网站日志,做出用户地域分布图》;网关层的地域风控看《OpenResty 网关怎么用 IP 归属地做地域风控》。