File size: 4,507 Bytes
bedc9c3
 
 
 
 
 
 
 
 
 
 
 
 
 
 
bcc58c0
bedc9c3
bcc58c0
bedc9c3
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
bcc58c0
bedc9c3
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
"""Node2Vec model for the edge_ML_expected_ge5 graph.

Default run is a smoke check (builds the model, runs a few optimizer steps).
Pass --epochs N to train, which writes embeddings to --out.
"""

import argparse
import time

import torch
from torch_geometric.data import Data
from torch_geometric.data.data import DataEdgeAttr, DataTensorAttr
from torch_geometric.data.storage import BaseStorage, EdgeStorage, GlobalStorage
from torch_geometric.nn import Node2Vec

from paths import EMB_PATH, GRAPH_PATH

GRAPH = GRAPH_PATH


def load_graph(path: str = GRAPH) -> Data:
    torch.serialization.add_safe_globals(
        [Data, DataEdgeAttr, DataTensorAttr, BaseStorage, EdgeStorage, GlobalStorage]
    )
    return torch.load(path, weights_only=True)


def build_model(data: Data, args: argparse.Namespace, device: torch.device) -> Node2Vec:
    return Node2Vec(
        data.edge_index,
        embedding_dim=args.embedding_dim,
        walk_length=args.walk_length,
        context_size=args.context_size,
        walks_per_node=args.walks_per_node,
        num_negative_samples=args.num_negative_samples,
        p=args.p,
        q=args.q,
        num_nodes=data.num_nodes,
        sparse=True,  # pairs with SparseAdam; the embedding table is the only param
    ).to(device)


def main() -> None:
    ap = argparse.ArgumentParser()
    ap.add_argument("--embedding-dim", type=int, default=128)
    ap.add_argument("--walk-length", type=int, default=20)
    ap.add_argument("--context-size", type=int, default=10)
    ap.add_argument("--walks-per-node", type=int, default=10)
    ap.add_argument("--num-negative-samples", type=int, default=1)
    ap.add_argument("--p", type=float, default=1.0, help="return parameter")
    ap.add_argument("--q", type=float, default=1.0, help="in-out parameter")
    ap.add_argument("--batch-size", type=int, default=128)
    ap.add_argument("--lr", type=float, default=0.01)
    ap.add_argument("--num-workers", type=int, default=4)
    ap.add_argument("--epochs", type=int, default=0, help="0 = smoke check only")
    ap.add_argument("--steps", type=int, default=5, help="steps for the smoke check")
    ap.add_argument("--out", default=EMB_PATH)
    args = ap.parse_args()

    device = torch.device("cuda" if torch.cuda.is_available() else "cpu")
    data = load_graph()
    model = build_model(data, args, device)

    print(f"graph      : {data.num_nodes:,} nodes, {data.edge_index.size(1) // 2:,} undirected edges")
    print(f"device     : {device}")
    print(f"model      : {model}")
    print(f"parameters : {sum(p.numel() for p in model.parameters()):,} "
          f"({data.num_nodes:,} x {args.embedding_dim})")
    print(f"walks      : length={args.walk_length} context={args.context_size} "
          f"per_node={args.walks_per_node} p={args.p} q={args.q}")

    loader = model.loader(batch_size=args.batch_size, shuffle=True,
                          num_workers=args.num_workers)
    optimizer = torch.optim.SparseAdam(list(model.parameters()), lr=args.lr)
    print(f"loader     : {len(loader):,} batches/epoch of {args.batch_size} seed nodes")

    def run_epoch(max_steps: int | None = None) -> float:
        model.train()
        total, n = 0.0, 0
        for i, (pos_rw, neg_rw) in enumerate(loader):
            optimizer.zero_grad()
            loss = model.loss(pos_rw.to(device), neg_rw.to(device))
            loss.backward()
            optimizer.step()
            total, n = total + loss.item(), n + 1
            if max_steps is not None and i + 1 >= max_steps:
                break
        return total / max(n, 1)

    if args.epochs == 0:
        t0 = time.perf_counter()
        loss = run_epoch(max_steps=args.steps)
        print(f"\nsmoke check: {args.steps} steps, mean loss {loss:.4f}, "
              f"{time.perf_counter() - t0:.1f}s")
        z = model()
        print(f"embeddings : {tuple(z.shape)} {z.dtype} on {z.device}")
        print("model built and training step verified; pass --epochs N to train")
        return

    for epoch in range(1, args.epochs + 1):
        t0 = time.perf_counter()
        loss = run_epoch()
        print(f"epoch {epoch:>3}/{args.epochs}  loss {loss:.4f}  "
              f"{time.perf_counter() - t0:.1f}s")

    model.eval()
    with torch.no_grad():
        z = model().cpu()
    torch.save({"embedding": z, "node_id": data.node_id,
                "args": vars(args)}, args.out)
    print(f"saved embeddings {tuple(z.shape)} -> {args.out}")


if __name__ == "__main__":
    main()