# Link Prediction for Heterogeneous graph

## 0. Enviroment setup

In [1]:
# !pip uninstall torch torchvision torchaudio --yes
# !pip install torch==2.2.1 torchvision==0.17.1 torchaudio==2.2.1 --index-url https://download.pytorch.org/whl/cu121
# !pip install lightning torch_geometric
# !pip install pyg_lib torch_scatter torch_sparse torch_cluster torch_spline_conv -f https://data.pyg.org/whl/torch-2.2.0+cu121.html
# !pip install wandb

In [2]:
import os
import shutil
import wandb
import torch
import torch.nn as nn
from torch.utils.data import Dataset
from torch.utils.data import DataLoader
from tqdm import tqdm

from torch_geometric.utils import negative_sampling
import torch_geometric.transforms as T
from torch_geometric.utils import train_test_split_edges
# from torch_geometric.loader import LinkNeighborLoader
# from torch_geometric.data.lightning import LightningLinkData

import lightning as L
from lightning.pytorch.callbacks import ModelCheckpoint
from lightning.pytorch.loggers import WandbLogger

In [3]:
from hive_analysis.models.link_prediction import *
from hive_analysis.dataloaders import hive_preprocessing

In [4]:
os.environ['CUDA_LAUNCH_BLOCKING'] = '1'

In [5]:
# Setup device agnostic code
device = "cuda" if torch.cuda.is_available() else "cpu"
device

'cuda'

## 1. Data Pre-processing

In [6]:
DATA_VERSION = 'final_v1'

In [7]:
# data = hive_preprocessing(
#     f'dataset/hive/{DATA_VERSION}/nodes_labelled.csv',
#     f'dataset/hive/{DATA_VERSION}/edges_labelled.csv',
#     to_undirected = True,
# )
# torch.save(data, f'dataset/hive/{DATA_VERSION}/hive.pt')
data = torch.load(f'dataset/hive/{DATA_VERSION}/hive.pt')
data

HeteroData(
  user={
    x=[18645, 5],
    node_label=<built-in method long of Tensor object at 0x7f62f97a30e0>,
  },
  comment={
    x=[125111, 1],
    node_label=<built-in method long of Tensor object at 0x7f62f97a3270>,
  },
  post={
    x=[13540, 1],
    node_label=<built-in method long of Tensor object at 0x7f62f97a34a0>,
  },
  (user, upvote, comment)={
    edge_index=[2, 423638],
    edge_attr=[423638, 1],
    y=[423638],
  },
  (user, upvote, post)={
    edge_index=[2, 554131],
    edge_attr=[554131, 1],
    y=[554131],
  },
  (user, write, comment)={
    edge_index=[2, 78696],
    edge_attr=[78696, 1],
    y=[78696],
  },
  (user, write, post)={
    edge_index=[2, 12958],
    edge_attr=[12958, 1],
    y=[12958],
  },
  (user, downvote, comment)={
    edge_index=[2, 6819],
    edge_attr=[6819, 1],
    y=[6819],
  },
  (user, downvote, post)={
    edge_index=[2, 2934],
    edge_attr=[2934, 1],
    y=[2934],
  },
  (comment, belong_to, comment)={
    edge_index=[2, 58911],
    ed

In [8]:
num_edges = len(data.edge_types)
edge_types = data.edge_types[:num_edges//2]

In [9]:
edge_types[0]

('user', 'upvote', 'comment')

In [10]:
rev_edge_types = data.edge_types[num_edges//2:]

In [11]:
transform = T.RandomLinkSplit(
    num_val=0.1,
    num_test=0.1,
    disjoint_train_ratio=0.3, 
    add_negative_train_samples=True,
    neg_sampling_ratio=2.0,
    edge_types=edge_types,
    rev_edge_types=rev_edge_types, 
)

train_data, val_data, test_data = transform(data)

## 2. Training

In [12]:
models = {
    # 'HGT': HGT, 
    # 'GATv2': GATv2, 
    # 'GraphSAGE': GraphSAGE, 
    # 'GAT':GAT
    'GraphConv': GraphConvNet,
}

In [13]:
models = { k: m(
    in_channels=-1,  
    out_channels=128,
    hidden_channels=[64, 128, 256, 256, 512], 
    metadata=data.metadata(), 
    edge_types=edge_types,
    rev_edge_types=rev_edge_types,
    # aggr_scheme='mean',
) for k, m in models.items()}

In [14]:
class GraphDataset(Dataset):
    def __init__(
        self,
        data,
        edge_types,
        key='edge_label',
    ):
        self.data = data
        self.edge_types = edge_types
        self.key =  key

    def __len__(self):
        return len(self.edge_types)

    def __getitem__(self, idx):
        return self.data, self.edge_types[idx], self.key
    
def collate_fn(input):
    data, edge_types, key = zip(*input)
    return data[0], edge_types, key[0]

train_loader = DataLoader(
    GraphDataset(train_data, edge_types),
    batch_size=1,
    shuffle=True,
    drop_last=False,
    pin_memory=True,
    num_workers=4,
    collate_fn=collate_fn,
)
val_loader = DataLoader(
    GraphDataset(val_data, edge_types),
    batch_size=len(edge_types),
    shuffle=False,
    drop_last=False,
    pin_memory=True,
    num_workers=4,
    collate_fn=collate_fn
)

In [15]:
# for edge_types, rev_edge_types in edges:
for mtype, model in models.items():
    log_dir = 'results/log/lp/' + mtype.lower()
    loss_checkpoint_dir = f'results/checkpoints/lp/{mtype.lower()}/loss'
    auc_checkpoint_dir = f'results/checkpoints/lp/{mtype.lower()}/roc_auc'
    acc_checkpoint_dir = f'results/checkpoints/lp/{mtype.lower()}/acc'

    os.makedirs(log_dir, exist_ok=True)
    os.makedirs(loss_checkpoint_dir, exist_ok=True)
    os.makedirs(auc_checkpoint_dir, exist_ok=True)
    os.makedirs(acc_checkpoint_dir, exist_ok=True)


    lr = 1e-3
    optim = torch.optim.Adam(model.parameters(), lr=lr)
    model.set_optimizer(optim)

    wandb_logger = WandbLogger(
        project="LinkPrediction_finalv1",
        log_model=True,
        save_dir=log_dir,
        name=mtype,
        entity='ssc_project'

    )

    loss_checkpoint_callback = ModelCheckpoint(
        monitor=f'val_loss',
        dirpath=loss_checkpoint_dir,
        filename='LinkPred-{epoch:02d}-{val_loss:.2f}',
        save_top_k=3,
        save_last=True,
        mode='min',
        every_n_epochs=1
    )
    roc_auc_checkpoint_callback = ModelCheckpoint(
        monitor=f'val_roc_auc',
        dirpath=auc_checkpoint_dir,
        filename='LinkPred-{epoch:02d}-{val_roc_auc:.2f}',
        save_top_k=3,
        save_last=True,
        mode='max',
        every_n_epochs=1
    )
    acc_checkpoint_callback = ModelCheckpoint(
        monitor=f'val_accuracy',
        dirpath=acc_checkpoint_dir,
        filename='LinkPred-{epoch:02d}-{val_accuracy:.2f}',
        save_top_k=3,
        save_last=True,
        mode='max',
        every_n_epochs=1
    )

    trainer = L.Trainer(
        max_epochs=500,
        check_val_every_n_epoch=10,
        callbacks=[
            loss_checkpoint_callback, 
            roc_auc_checkpoint_callback,
            acc_checkpoint_callback,
        ],
        logger=wandb_logger,
        log_every_n_steps=4
    )


    trainer.fit(model, train_loader, val_loader)
    wandb.finish()

GPU available: True (cuda), used: True
TPU available: False, using: 0 TPU cores
IPU available: False, using: 0 IPUs
HPU available: False, using: 0 HPUs
You are using a CUDA device ('NVIDIA RTX A6000') that has Tensor Cores. To properly utilize them, you should set `torch.set_float32_matmul_precision('medium' | 'high')` which will trade-off precision for performance. For more details, read https://pytorch.org/docs/stable/generated/torch.set_float32_matmul_precision.html#torch.set_float32_matmul_precision
[34m[1mwandb[0m: Currently logged in as: [33mhontrn9122[0m ([33mssc_project[0m). Use [1m`wandb login --relogin`[0m to force relogin


/usr/local/lib/python3.9/dist-packages/lightning/pytorch/callbacks/model_checkpoint.py:653: Checkpoint directory /notebooks/results/checkpoints/lp/graphconv/loss exists and is not empty.
/usr/local/lib/python3.9/dist-packages/lightning/pytorch/callbacks/model_checkpoint.py:653: Checkpoint directory /notebooks/results/checkpoints/lp/graphconv/roc_auc exists and is not empty.
/usr/local/lib/python3.9/dist-packages/lightning/pytorch/callbacks/model_checkpoint.py:653: Checkpoint directory /notebooks/results/checkpoints/lp/graphconv/acc exists and is not empty.
LOCAL_RANK: 0 - CUDA_VISIBLE_DEVICES: [0]
/usr/local/lib/python3.9/dist-packages/lightning/pytorch/utilities/model_summary/model_summary.py:454: A layer with UninitializedParameter was found. Thus, the total number of parameters detected may be inaccurate.

  | Name    | Type              | Params
----------------------------------------------
0 | encoder | ModuleDict        | 77.8 M
1 | crit    | BCEWithLogitsLoss | 0     
---------

Sanity Checking: |          | 0/? [00:00<?, ?it/s]

Training: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

Validation: |          | 0/? [00:00<?, ?it/s]

`Trainer.fit` stopped: `max_epochs=500` reached.


0,1
epoch,▁▁▁▁▂▂▂▂▂▃▃▃▃▃▃▄▄▄▄▄▅▅▅▅▅▅▆▆▆▆▆▆▇▇▇▇▇███
train_loss,▆▅▃▃█▄▄▄▄▃▃▃▂▁▁▃▁▆▂▆▃▃▁▁▂▄▂▂▁▂▂▂▁▄▂▁▁▂▁▁
trainer/global_step,▁▁▁▁▂▂▂▂▂▃▃▃▃▃▃▄▄▄▄▄▄▅▅▅▅▅▆▆▆▆▆▆▇▇▇▇▇███
val_accuracy,▁▃▅▆▆▇▆▆▆▇▇▇▇█▇█████████████████████████
val_f1,▁▃▅▆▆▇▆▆▆▇▇▇▇█▇█████████████████████████
val_loss,█▆▄▄▃▂▃▂▃▂▂▂▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁
val_precision,▁▂▄▄▅▆▄▅▄▆▅▅▆▇▆▇▇▆▇▇▇██▇▇▇▇▇█▇██▇███████
val_recall,▁▅▅▆▇▆█▇█▇██▇▇█▇▇▇▇▇▇▆▆▇▇▇▇▆▆▆▆▆▆▆▆▆▆▆▆▆
val_roc_auc,▁▄▅▆▇▇▇▇▇▇███████████▇████████▇▇▇▇▇▇▇▇▇▇

0,1
epoch,499.0
train_loss,0.37563
trainer/global_step,3999.0
val_accuracy,0.86092
val_f1,0.86092
val_loss,0.4563
val_precision,0.76702
val_recall,0.83697
val_roc_auc,0.85493


In [16]:
print('done')

done
