Skip to content

物联网设备接入了哪些省?用 Go 给设备连接日志打 IP 地域标签(含 NAT 出口的坑) ​

一个做智能水表的团队问过我一个很朴素的问题:设备装出去两万四千台,究竟装在哪些省?

他们的平台里,设备上报只带设备号和业务数据,没有定位模块(加 GPS 是钱,还得担心信号)。但每台设备连平台时,服务端都知道它从哪个 IP 来的。于是问题变成:把 IP 换成省份,再按设备数统计一遍。

听着简单,真正做的时候有三个坑等着:设备日志里的 IP 是网关地址不是设备地址;几百台设备共用一个 NAT 出口 IP,按 IP 计数会得出荒谬的结论;两万四千台设备一个个查接口,速率这道坎绕不过去。

这篇给一个 Go 单文件工具,读设备连接日志,输出分省设备数、出口运营商分布,外加一份"需要人工确认"的清单。

一、场景问题 ​

设备侧的统计和网站侧不太一样,差别主要在数据形态上。

网站日志一行一个访客,设备日志一行一台设备连一次。 同一台设备一天可能上线十几次,同一个 IP 底下可能挂着几百台设备。所以"按 IP 统计"和"按设备统计"是两个结果,前面那个会严重失真 —— 我们实测的样本里,一个出口 IP 底下挂了 3 台设备就足以让省份排名翻转。

真实 IP 不一定在日志字段里。 平台前面有负载均衡或者反向代理时,服务端看到的连接来源是代理的地址,真实客户端 IP 在 X-Forwarded-For 这类头部里;MQTT 走 TCP 的话没有 HTTP 头可用,得靠 PROXY protocol 或者 EMQX 的规则引擎把 peername 取出来落进日志。这一步不在本文范围内,但打标签之前一定要确认自己手里那个 IP 是设备真实的,否则后面全省份统计都是错的。

没人盯着看。 设备分布是日报/周报级的东西,偶尔跑一次就行,所以不需要常驻服务,一个能塞进 crontab 的命令行工具最合适。

二、方案设计 ​

程序分四段,每段解决一个问题:

  1. 去重:日志里先按 IP 归组,记录每个 IP 对应哪些设备号。查接口的次数降到"唯一 IP 数",而不是"日志行数"。
  2. 限速:用一个令牌桶 goroutine 统一发令牌,worker 只管等令牌再发请求。限速集中在一处,加减并发都不用改业务代码。默认 1 次/秒 —— 免费版额度是 60 次/分钟/IP,留足余量。
  3. 缓存:进程内 map 加互斥锁,同时把结果逐行追加到 ip_cache.jsonl。第二次跑同类日志基本零请求(实测 4.2 秒 → 271 毫秒)。
  4. 归类与汇总:把结果分成三类 —— 国内(按 prov/city 统计)、境外(按 country 统计)、非公网或查询失败(单独列出来人工看)。按设备台数计权重,不是按 IP 数。

字段上只需要 country、prov、city、isp、big_area,免费版接口都给。

三、Go 实现 ​

完整代码(iotstat.go,纯标准库):

go
// 给物联网平台的设备连接日志打 IP 地域标签,输出分省设备统计
// go build -o iotstat iotstat.go
// ./iotstat -in devices.tsv -qps 2 -out devices_tagged.tsv
package main

import (
	"bufio"
	"encoding/json"
	"flag"
	"fmt"
	"net"
	"net/http"
	"os"
	"sort"
	"strings"
	"sync"
	"time"
)

const apiBase = "https://ip9.com.cn/get"

type geo struct {
	IP      string `json:"ip"`
	Country string `json:"country"`
	Prov    string `json:"prov"`
	City    string `json:"city"`
	ISP     string `json:"isp"`
	BigArea string `json:"big_area"`
}

type apiResp struct {
	Ret int `json:"ret"`
	// 非法 IP 时接口返回 ret=400 且 data 是空数组,用 RawMessage 接住再二次解析
	Data json.RawMessage `json:"data"`
}

type result struct {
	IP   string
	Geo  geo
	Kind string // cn=国内 oversea=境外 internal=非公网 error=查询失败
	Note string
}

var (
	inFile    = flag.String("in", "devices.tsv", "设备连接日志(device_id<TAB>ip)")
	outFile   = flag.String("out", "", "打标结果写 TSV(留空不写)")
	qps       = flag.Float64("qps", 1.0, "请求速率上限,免费版 60 次/分钟,建议 1")
	workers   = flag.Int("workers", 4, "并发 worker 数")
	cacheFile = flag.String("cache", "ip_cache.jsonl", "IP 归属地缓存文件")
)

type cache struct {
	mu sync.Mutex
	m  map[string]geo
}

func newCache() *cache { return &cache{m: make(map[string]geo)} }

func (c *cache) get(ip string) (geo, bool) {
	c.mu.Lock()
	defer c.mu.Unlock()
	g, ok := c.m[ip]
	return g, ok
}

func (c *cache) put(ip string, g geo) {
	c.mu.Lock()
	c.m[ip] = g
	c.mu.Unlock()
	f, err := os.OpenFile(*cacheFile, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o644)
	if err != nil {
		return
	}
	defer f.Close()
	json.NewEncoder(f).Encode(g)
}

func loadCache() *cache {
	c := newCache()
	f, err := os.Open(*cacheFile)
	if err != nil {
		return c
	}
	defer f.Close()
	sc := bufio.NewScanner(f)
	for sc.Scan() {
		var g geo
		if json.Unmarshal(sc.Bytes(), &g) == nil && g.IP != "" {
			c.m[g.IP] = g
		}
	}
	fmt.Printf("缓存载入 %d 条\n", len(c.m))
	return c
}

func queryIP(client *http.Client, ip string) (geo, error) {
	req, err := http.NewRequest("GET", apiBase+"?ip="+ip, nil)
	if err != nil {
		return geo{}, err
	}
	req.Header.Set("User-Agent", "iotstat/1.0")
	resp, err := client.Do(req)
	if err != nil {
		return geo{}, err
	}
	defer resp.Body.Close()
	if resp.StatusCode == 429 {
		return geo{}, fmt.Errorf("触发限速 ret=429")
	}
	var r apiResp
	if err := json.NewDecoder(resp.Body).Decode(&r); err != nil {
		return geo{}, err
	}
	if resp.StatusCode != 200 || r.Ret != 200 {
		return geo{}, fmt.Errorf("HTTP %d ret=%d", resp.StatusCode, r.Ret)
	}
	var g geo
	if err := json.Unmarshal(r.Data, &g); err != nil {
		return geo{}, fmt.Errorf("data 解析失败:%w", err)
	}
	return g, nil
}

func isPrivate(ipStr string) bool {
	ip := net.ParseIP(ipStr)
	if ip == nil {
		return false
	}
	for _, cidr := range []string{"10.0.0.0/8", "172.16.0.0/12", "192.168.0.0/16",
		"127.0.0.0/8", "169.254.0.0/16", "fc00::/7", "fe80::/10"} {
		_, n, _ := net.ParseCIDR(cidr)
		if n.Contains(ip) {
			return true
		}
	}
	return false
}

func readDevices(path string) ([]string, []string, error) {
	f, err := os.Open(path)
	if err != nil {
		return nil, nil, err
	}
	defer f.Close()
	var ids, ips []string
	sc := bufio.NewScanner(f)
	for sc.Scan() {
		line := strings.TrimSpace(sc.Text())
		if line == "" || strings.HasPrefix(line, "#") {
			continue
		}
		parts := strings.Fields(line)
		if len(parts) < 2 {
			continue
		}
		ids = append(ids, parts[0])
		ips = append(ips, parts[1])
	}
	return ids, ips, sc.Err()
}

func main() {
	flag.Parse()
	ids, ips, err := readDevices(*inFile)
	if err != nil {
		fmt.Println("读取设备日志失败:", err)
		os.Exit(1)
	}

	// 同一个出口 IP 只查一次,记住它底下挂了哪些设备
	uniq := map[string][]string{}
	for i, ip := range ips {
		uniq[ip] = append(uniq[ip], ids[i])
	}
	fmt.Printf("连接记录 %d 条,去重后 %d 个 IP,速率上限 %.1f req/s\n", len(ips), len(uniq), *qps)

	// 令牌桶:把速率控制集中在一处
	tokens := make(chan struct{}, 1)
	go func() {
		interval := time.Duration(float64(time.Second) / *qps)
		for {
			tokens <- struct{}{}
			time.Sleep(interval)
		}
	}()

	client := &http.Client{Timeout: 6 * time.Second}
	ipCache := loadCache()

	ipList := make([]string, 0, len(uniq))
	for ip := range uniq {
		ipList = append(ipList, ip)
	}
	sort.Strings(ipList)

	results := make([]result, len(ipList))
	var wg sync.WaitGroup
	sem := make(chan struct{}, *workers)
	start := time.Now()
	for i, ip := range ipList {
		wg.Add(1)
		sem <- struct{}{}
		go func(i int, ip string) {
			defer wg.Done()
			defer func() { <-sem }()
			if g, ok := ipCache.get(ip); ok {
				results[i] = classify(ip, g, "缓存命中")
				return
			}
			<-tokens
			g, err := queryIP(client, ip)
			if err != nil {
				if isPrivate(ip) {
					results[i] = result{IP: ip, Kind: "internal", Note: "内网/链路本地地址,无地域意义"}
				} else {
					results[i] = result{IP: ip, Kind: "error", Note: err.Error()}
				}
				return
			}
			ipCache.put(ip, g)
			results[i] = classify(ip, g, "")
		}(i, ip)
	}
	wg.Wait()
	fmt.Printf("\n查询完成,耗时 %s(%d 个唯一 IP)\n",
		time.Since(start).Round(time.Millisecond), len(ipList))

	// —— 汇总:权重是设备台数,不是 IP 数 ——
	provCount, ispCount := map[string]int{}, map[string]int{}
	var problems []result
	for _, r := range results {
		n := len(uniq[r.IP]) // 这个出口 IP 底下挂了 n 台设备
		switch r.Kind {
		case "cn":
			provCount[r.Geo.Prov+" "+r.Geo.City] += n
			ispCount[r.Geo.ISP] += n
		case "oversea":
			provCount[r.Geo.Country] += n
		default:
			problems = append(problems, r)
		}
	}

	type kv struct {
		K string
		V int
	}
	rank := func(m map[string]int) []kv {
		out := make([]kv, 0, len(m))
		for k, v := range m {
			out = append(out, kv{k, v})
		}
		sort.Slice(out, func(i, j int) bool { return out[i].V > out[j].V })
		return out
	}
	fmt.Println("\n—— 设备属地 TOP ——")
	for _, r := range rank(provCount) {
		fmt.Printf("%-16s %d 台\n", r.K, r.V)
	}
	fmt.Println("\n—— 出口运营商 TOP ——")
	for _, r := range rank(ispCount) {
		fmt.Printf("%-30s %d 台\n", r.K, r.V)
	}
	fmt.Println("\n—— 需要人工看的记录 ——")
	for _, r := range problems {
		fmt.Printf("%-22s %-9s %s\n", r.IP, r.Kind, r.Note)
	}
	fmt.Println("\n—— 共用出口 IP(NAT 后的设备)——")
	for ip, devs := range uniq {
		if len(devs) > 1 {
			fmt.Printf("%-24s %d 台  %s\n", ip, len(devs), strings.Join(devs, ","))
		}
	}

	if *outFile != "" {
		f, _ := os.Create(*outFile)
		defer f.Close()
		w := bufio.NewWriter(f)
		defer w.Flush()
		for _, r := range results {
			fmt.Fprintf(w, "%s\t%s\t%s\t%s\t%s\t%s\n",
				r.IP, r.Kind, r.Geo.Country, r.Geo.Prov, r.Geo.City, r.Geo.ISP)
		}
		fmt.Println("\n已写出", *outFile)
	}
}

func classify(ip string, g geo, note string) result {
	switch {
	case g.Country == "中国":
		return result{IP: ip, Geo: g, Kind: "cn", Note: note}
	case g.Country != "" && g.Country != "保留":
		return result{IP: ip, Geo: g, Kind: "oversea", Note: note}
	default:
		return result{IP: ip, Geo: g, Kind: "internal", Note: "非公网地址:" + g.ISP}
	}
}

日志格式就是两列,设备号和 IP:

DEV-1001	114.114.114.114
DEV-1002	223.5.5.5
DEV-1007	192.168.1.1

sort.Strings 排一下再并发,是为了让输出顺序稳定、和缓存文件的顺序一致 —— 排查问题时这点很省事。

四、实测输出 ​

14 条连接记录、9 个唯一 IP,速率设成 2 次/秒:

连接记录 14 条,去重后 9 个 IP,速率上限 2.0 req/s

查询完成,耗时 4.236s(9 个唯一 IP)

—— 设备属地 TOP ——
江苏 南京          6 台
浙江 杭州          3 台
广东 广州          1 台
捷克             1 台

—— 出口运营商 TOP ——
中国电信                         5 台
AliDNS/DoH/DoT/阿里云           3 台
114DNS                       2 台

—— 需要人工看的记录 ——
192.168.1.1            internal  非公网地址:内网地址
203.0.113.9            internal  非公网地址:文档地址
999.1.1.1              error     HTTP 400 ret=400

—— 共用出口 IP(NAT 后的设备)——
223.5.5.5                3 台  DEV-1002,DEV-1010,DEV-1014
58.213.1.1               2 台  DEV-1003,DEV-1012
114.114.114.114          2 台  DEV-1001,DEV-1006
180.101.49.12            2 台  DEV-1004,DEV-1013

几个细节值得说:

192.168.1.1 和 203.0.113.9 不算"查询失败"。 接口对它们返回 HTTP 200,country=保留,isp 分别是"内网地址"和"文档地址"。前者出现在日志里说明你拿到的是内网地址(链路或配置有问题),后者说明测试数据混进了生产日志。这两类要和"接口报错"分开列,否则排查方向就错了。

999.1.1.1 是 HTTP 400 + ret=400。 非法 IP 接口在 HTTP 层和 body 层都标了错,两层都要判:只看 HTTP 状态码会漏掉 body 里的 ret=429(限速),只看 ret 会漏掉 HTTP 层的异常。

第二次跑走缓存的差距很夸张。 同样的日志再跑一遍:

缓存载入 8 条
查询完成,耗时 271ms(9 个唯一 IP)

4.2 秒变 271 毫秒,因为 8 个成功的 IP 都命中了缓存,剩下那个非法的没写缓存(失败不缓存,避免把一次网络抖动固化下来)。

五、真上线要补的几件事 ​

先确认你手里是设备真实 IP。 平台前面有 SLB/NGINX 时,取 X-Forwarded-For 里第一个非内网地址;MQTT 走 TCP 的场景要靠 PROXY protocol 拿到对端地址。这一步做错,后面的省份分布会整齐地全部指向你的机房所在城市 —— 这个错误现象很好认:所有设备都在同一个城市,基本都是取到了代理 IP。

NAT 后的设备要单独说明口径。 一个小区、一个工厂、一个运营商大内网出口后面挂着几百台设备,它们共享一个 IP,归属地只有一个。所以报表里容易出现"江苏南京 6 台"这种看着很干净的分布,实际含义是"这 6 台设备的流量从江苏南京的出口出来"。做区域盘点时,最好把"共用出口"这一栏一起给业务方看,否则容易被理解成"设备就装在南京"。

量级上去之后,策略要换。 两万四千台设备,去重后假设一万五千个唯一 IP,免费版 60 次/分钟要跑 4 个多小时 —— 加缓存、每天只处理新增 IP 也能接受。但如果是一百万台量级的平台,就得考虑 VIP 版(18 万次/分钟)或者干脆用离线库在本地跑,热数据走在线接口、冷数据走离线表,两边对不上时以在线为准。

IPv6 设备越来越多,别在日志格式上被卡住。 采集阶段按字符串存 IP 最省事,别自作聪明转成整数 —— IPv4 的整数表达(long_ip)接口会一起返回,IPv6 没有这个概念,转来转去反而丢信息。

合规上留个心。 设备 IP 在个人信息保护语境下属于设备网络标识,做统计可以,别把它跟设备使用者的身份信息绑死存进业务库。报表只保留下钻到省市级的聚合结果,原始 IP 用一段时间后清理,这条在给客户交付数据时是最容易被问到的。

统计口径确认之后,这件事的代码量其实很小:一个限速器、一层缓存、一次分组。接口 GET https://ip9.com.cn/get?ip=<ip> 免费版 60 次/分钟,返回的 prov、city、isp、big_area 就够做设备分布报表;要区县或者 ip_type(区分家庭宽带和 IDC 机房)就得看 VIP 版。完整字段和更新说明在官网 https://www.ip9.com.cn 上。工具可以照搬给别的场景用:CDN 节点的客户端分布、门店的访客来源核查、渠道设备的落地区域核验,都是同一个套路。