In [2]:
import torch
import sys

sys.path.append("/workspace/kbqa/")  # go to parent dir
device = torch.device("cuda" if torch.cuda.is_available() else "cpu")
device

device(type='cpu')

In [3]:
from tqdm import tqdm
import jsonlines
import networkx as nx
import pandas as pd
import numpy as np
import torch
import json
import pandas as pd
from tqdm import tqdm
from torch.utils.data import Dataset
from transformers import TrainingArguments, Trainer
from transformers.models.graphormer.collating_graphormer import GraphormerDataCollator
from transformers import GraphormerForGraphClassification
from transformers.models.graphormer.collating_graphormer import algos_graphormer

2023-08-08 08:44:54.778166: I tensorflow/core/platform/cpu_feature_guard.cc:193] This TensorFlow binary is optimized with oneAPI Deep Neural Network Library (oneDNN) to use the following CPU instructions in performance-critical operations:  AVX2 FMA
To enable them in other operations, rebuild TensorFlow with the appropriate compiler flags.
2023-08-08 08:44:55.042228: E tensorflow/stream_executor/cuda/cuda_blas.cc:2981] Unable to register cuBLAS factory: Attempting to register factory for plugin cuBLAS when one has already been registered
2023-08-08 08:44:55.661800: W tensorflow/stream_executor/platform/default/dso_loader.cc:64] Could not load dynamic library 'libnvinfer.so.7'; dlerror: libnvinfer.so.7: cannot open shared object file: No such file or directory; LD_LIBRARY_PATH: /usr/local/nvidia/lib:/usr/local/nvidia/lib64
2023-08-08 08:44:55.661885: W tensorflow/stream_executor/platform/default/dso_loader.cc:64] Could not load dynamic library 'libnvinfer_plugin.so.7'; dlerror: libnvinf

In [10]:
dataset_type = 't5-xl-ssm'
train_bs = 64
eval_bs = 64
data_prep = False
push_to_hub = True
model_weights = "/workspace/storage/subgraphs_reranking_results/t5-xl-ssm/results/clefourrier/graphormer-base-pcqm4mv2_mse/checkpoint-119500"
model_name = "clefourrier/graphormer-base-pcqm4mv2"
num_epochs = 50
model_save_name = f"{model_name}_mse" if not model_weights else model_weights.split('/')[-1]

In [5]:
def read_jsonl(path):
    jsonl_reader = jsonlines.open(path)
    jsonl_reader_list = list(jsonl_reader)
    df = []
    for line in tqdm(jsonl_reader_list):
        df.append(line)
    df = pd.DataFrame(df)
    return df

def preprocess_df(df):
    # turn list of entities into string
    df["answerEntity"] = df["answerEntity"].apply(lambda x: ", ".join(x))
    df["questionEntity"] = df["questionEntity"].apply(lambda x: ", ".join(x))
    df["groundTruthAnswerEntity"] = df["groundTruthAnswerEntity"].apply(
        lambda x: ", ".join(x)
    )
    df["correct"] = df.apply(
        lambda x: x["answerEntity"] in x["groundTruthAnswerEntity"], axis=1
    )
    
    return df

In [6]:
if data_prep:
    train_df = read_jsonl(
        f"/workspace/storage/new_subgraph_dataset/{dataset_type}/mintaka_train_labeled.jsonl"
    )
    train_df = preprocess_df(train_df)
    val_df = read_jsonl(
        f"/workspace/storage/new_subgraph_dataset/{dataset_type}/mintaka_validation_labeled.jsonl"
    )
    val_df = preprocess_df(val_df)
    test_df = read_jsonl(
        f"/workspace/storage/new_subgraph_dataset/{dataset_type}/mintaka_test_labeled.jsonl"
    )
    test_df = preprocess_df(test_df)

#### Transforming the data

In [7]:
def transform_graph(graph_data, answer_entity, ground_truth_entity):
    # Create an empty dictionary to store the transformed graph
    transformed_graph = {}

    # Extract 'nodes' and 'links' from the graph_data
    nodes = graph_data['nodes']
    links = graph_data['links']

    # Calculate num_nodes
    num_nodes = len(nodes)

    # Calculate edge_index
    edge_index = [[link['source'], link['target']] for link in links]
    edge_index = list(zip(*edge_index))

    # Check if "answerEntity" matches with "groundTruthAnswerEntity" to get the label (y)
    y = 1.0 if answer_entity in ground_truth_entity else 0.0

    # Calculate node_feat based on 'type' key
    node_feat = []
    for node in nodes:
        if node['type'] == 'INTERNAL':
            node_feat.append([1])
        elif node['type'] == 'ANSWER_CANDIDATE_ENTITY':
            node_feat.append([2])
        elif node['type'] == 'QUESTIONS_ENTITY':
            node_feat.append([3])
    
    # Store the calculated values in the transformed_graph dictionary
    transformed_graph['edge_index'] = edge_index
    transformed_graph['num_nodes'] = num_nodes
    transformed_graph['y'] = [y]
    transformed_graph['node_feat'] = node_feat
    transformed_graph['edge_attr'] = [[0]]

    return transformed_graph


In [8]:
def create_adjacency_matrix(edge_list):
    # Find the maximum node ID in the edge_list
    max_node_id = max(max(edge_list[0]), max(edge_list[1]))

    # Initialize an empty adjacency matrix with zeros
    adjacency_matrix = np.zeros((max_node_id+1, max_node_id+1), dtype=np.int32)  

    # Add edges to the adjacency matrix
    for src, dest in zip(edge_list[0], edge_list[1]):
        adjacency_matrix[src, dest] = 1  
    

    return adjacency_matrix

In [9]:
def preprocess(item):
    """Convert to the required format for Graphormer"""
    attn_edge_type = None  # Initialize outside the loop

    # Calculate adjacency matrix
    adj = create_adjacency_matrix(item["edge_index"])

    shortest_path_result, path = algos_graphormer.floyd_warshall(adj)

    try:
        # Calculate max_dist and input_edges if the function call succeeds
        shortest_path_result, path = algos_graphormer.floyd_warshall(adj)
        max_dist = np.amax(shortest_path_result)
        attn_edge_type = np.zeros((item["num_nodes"], item["num_nodes"], len(item['edge_attr'])), dtype=np.int64)
        input_edges = algos_graphormer.gen_edge_input(max_dist, path, attn_edge_type)
    except:
        # If the function call fails, handle the exception
        max_dist = 0
        attn_edge_type = None
        input_edges = np.zeros((item["num_nodes"], item["num_nodes"], max_dist, len(item['edge_attr'])), dtype=np.int64)
        shortest_path_result = None

    if attn_edge_type is None:
        # Initialize attn_edge_type here if it hasn't been initialized already
        attn_edge_type = np.zeros((item["num_nodes"], item["num_nodes"], len(item['edge_attr'])), dtype=np.int64)

    # Set values for all the keys
    processed_item = {
        "edge_index": np.array(item["edge_index"]),
        "num_nodes": item["num_nodes"],
        "y": item["y"],
        "node_feat": np.array(item["node_feat"]),
        "input_nodes": np.array(item["node_feat"]),  # Use node_feat as input_nodes if node_feat is the feature representation
        "edge_attr": np.array(item["edge_attr"]),
        "attn_bias": np.zeros((item["num_nodes"] + 1, item["num_nodes"] + 1), dtype=np.single),
        "attn_edge_type": attn_edge_type,
        "spatial_pos": shortest_path_result.astype(np.int64) + 1,
        "in_degree": np.sum(adj, axis=1).reshape(-1) + 1,
        "out_degree": np.sum(adj, axis=1).reshape(-1) + 1,  # for undirected graph
        "input_edges": input_edges + 1,
        "labels": item.get("labels", item["y"]),  # Assuming "labels" key may or may not exist in the input data
    }

    return processed_item

In [10]:
def transform_data(df, save_path):
    transformed_graph_dicts = []
    for _, row in tqdm(df.iterrows(), total=len(df), desc="Transforming graphs"):
        try:
            curr_dict = {}
            graph_data = row['graph']
            curr_dict['original_graph'] = graph_data

            transformed_graph = transform_graph(graph_data, row['answerEntity'], row['groundTruthAnswerEntity'])
            if len(transformed_graph["edge_index"][0]) or len(transformed_graph["edge_index"][1]) > 1:
                curr_dict['question'] = row['question']
                curr_dict['answerEntity'] = row['answerEntity']
                curr_dict['groundTruthAnswerEntity'] = row['groundTruthAnswerEntity']
                curr_dict['correct'] = row['correct']
                curr_dict['transformed_graph'] = transformed_graph
                transformed_graph_dicts.append(curr_dict)
        except:
            continue 

            
    with open(save_path, 'w+') as file:
        for transformed_graph in transformed_graph_dicts:
            file.write(json.dumps(transformed_graph) + '\n')

In [11]:
# filter out yesno and count questions
if data_prep:
    train_df = pd.concat([train_df, val_df])
    train_df = train_df[(train_df['complexityType'] != 'yesno') & (train_df['complexityType'] != 'count')] 
    test_df = test_df[(test_df['complexityType'] != 'yesno') & (test_df['complexityType'] != 'count')] 

In [12]:
train_trans_path = f'/workspace/storage/new_subgraph_dataset/{dataset_type}/graph_class/transformed_graphs_train_1.jsonl'
test_trans_path = f'/workspace/storage/new_subgraph_dataset/{dataset_type}/graph_class/transformed_graphs_test_1.jsonl'

In [13]:
if data_prep:
    transform_data(test_df, test_trans_path)
    transform_data(train_df, train_trans_path)

### Preparing the data

In [14]:
from ast import literal_eval
from unidecode import unidecode
def try_literal_eval(s):
    try:
        return literal_eval(s)
    except ValueError:
        print('yo')
        return s


class CustomGraphDataset(Dataset):
    def __init__(self, file_path):
        self.data = []
        with open(file_path, 'r') as file:
            for line in file:
                graph_dicts = json.loads(line)
                preproc_graph = preprocess(graph_dicts['transformed_graph'])
                
                if preproc_graph['input_edges'].shape[2] != 0:
                    self.data.append(preproc_graph)
    
    def __len__(self):
        return len(self.data)

    def __getitem__(self, idx):
        return self.data[idx]

# Load your custom training and test datasets
train_dataset = CustomGraphDataset(train_trans_path)
test_dataset = CustomGraphDataset(test_trans_path)

In [15]:
len(train_dataset)

79752

#### Training

In [16]:
import numpy as np
import evaluate


threshold = 0.5
metric_classifier = evaluate.combine(["accuracy", "f1", "precision", "recall", "hyperml/balanced_accuracy",])
metric_regression = evaluate.combine(["mae"])


def compute_metrics(eval_pred):
    predictions, labels = eval_pred
    predictions = predictions[0]
    results = metric_regression.compute(predictions=predictions, references=labels)

    predictions = predictions > threshold
    results.update(
        metric_classifier.compute(predictions=predictions, references=labels)
    )
    return results

In [6]:
model_weights

'/workspace/storage/subgraphs_reranking_results/t5-large-ssm/results/clefourrier/graphormer-base-pcqm4mv2_mse/checkpoint-62000'

In [11]:
if  model_weights: # evaluating previous trained model weights
    model = GraphormerForGraphClassification.from_pretrained(
    model_weights,
    num_classes=1,
    ignore_mismatched_sizes=True,)
    
    # push this version to the hub
    if push_to_hub:
        model.push_to_hub(commit_message='previous trained best checkpoint', repo_id=f'hle2000/graphsormer_subgraphs_reranking_{dataset_type}')
else: # training from scratch
    model = GraphormerForGraphClassification.from_pretrained(
    model_name,
    num_classes=1,
    ignore_mismatched_sizes=True,)

pytorch_model.bin:   0%|          | 0.00/191M [00:00<?, ?B/s]

Upload 1 LFS files:   0%|          | 0/1 [00:00<?, ?it/s]

In [18]:
from torch.utils.data.sampler import WeightedRandomSampler
import numpy as np

class CustomTrainer(Trainer):  
    def get_labels(self):
        labels = []
        for i in self.train_dataset:
            labels.append(int(i["y"][0]))
        return labels

    def _get_train_sampler(self) -> torch.utils.data.Sampler:
        labels = self.get_labels()
        return self.create_sampler(labels)
      
    def create_sampler(self, target):
        class_sample_count = np.array(
            [len(np.where(target == t)[0]) for t in np.unique(target)]
        )
        weight = 1.0 / class_sample_count
        samples_weight = np.array([weight[t] for t in target])

        samples_weight = torch.from_numpy(samples_weight)
        samples_weigth = samples_weight.double()
        sampler = WeightedRandomSampler(samples_weight, len(samples_weight))

        return sampler

In [19]:
# Specifiy the arguments for the trainer
training_args = TrainingArguments(
    output_dir=f"/workspace/storage/subgraphs_reranking_results/{dataset_type}/results/{model_save_name}",  # output directory
    num_train_epochs=num_epochs,  # total number of training epochs
    per_device_train_batch_size=train_bs,  # batch size per device during training
    per_device_eval_batch_size=eval_bs,  # batch size for evaluation
    warmup_steps=500,  # number of warmup steps for learning rate scheduler
    weight_decay=0.01,  # strength of weight decay
    logging_dir=f"/workspace/storage/subgraphs_reranking_results/{dataset_type}/logs/{model_save_name}",  # directory for storing logs
    load_best_model_at_end=True,  # load the best model when finished training (default metric is loss)
    metric_for_best_model="balanced_accuracy",  # select the base metrics
    logging_steps=500,  # log & save weights each logging_steps
    save_steps=500,
    evaluation_strategy="steps",  # evaluate each `logging_steps`
    report_to='wandb',
)

In [20]:
# Initialize the data collator
data_collator = GraphormerDataCollator()
# Initialize the Trainer
trainer = CustomTrainer(
    model=model,
    args=training_args,
    train_dataset=train_dataset,
    eval_dataset=test_dataset,
    data_collator=data_collator,
    compute_metrics=compute_metrics,  # the callback that computes metrics of interest
)

In [21]:
if not model_weights: # training
    train_results = trainer.train()
    trainer.save_model(f"/workspace/storage/subgraphs_reranking_results/{dataset_type}/results/{model_save_name}/best_checkpoint")
    if push_to_hub:
        trainer.push_to_hub(commit_message='best checkpoint')

### Eval

In [22]:
evaluate_res = trainer.evaluate()
evaluate_res

Failed to detect the name of this notebook, you can set it manually with the WANDB_NOTEBOOK_NAME environment variable to enable code saving.
[34m[1mwandb[0m: Currently logged in as: [33mhle2000[0m. Use [1m`wandb login --relogin`[0m to force relogin


{'eval_loss': 0.20189997553825378,
 'eval_mae': 0.37012999221773013,
 'eval_accuracy': 0.6751517779705117,
 'eval_f1': 0.33501997336884154,
 'eval_precision': 0.22059855038578444,
 'eval_recall': 0.6960531169310218,
 'eval_balanced_accuracy': 0.6842101547110266,
 'eval_runtime': 46.0724,
 'eval_samples_per_second': 500.517,
 'eval_steps_per_second': 7.835}

#### Re-ranking

In [23]:
test_df = read_jsonl(test_trans_path)

100%|██████████| 23070/23070 [00:00<00:00, 2472722.92it/s]


In [24]:
test_df.head()

Unnamed: 0,original_graph,question,answerEntity,groundTruthAnswerEntity,correct,transformed_graph
0,"{'directed': True, 'multigraph': False, 'graph...",What man was a famous American author and also...,Q893594,Q7245,False,"{'edge_index': [[0, 1, 2, 2, 3, 3], [0, 0, 0, ..."
1,"{'directed': True, 'multigraph': False, 'graph...",What man was a famous American author and also...,Q102513,Q7245,False,"{'edge_index': [[1, 1, 2, 2, 3, 4, 4], [0, 1, ..."
2,"{'directed': True, 'multigraph': False, 'graph...",What man was a famous American author and also...,Q7245,Q7245,True,"{'edge_index': [[1, 3, 3], [0, 0, 2]], 'num_no..."
3,"{'directed': True, 'multigraph': False, 'graph...",What man was a famous American author and also...,Q34652890,Q7245,False,"{'edge_index': [[0, 1, 1, 2, 3, 3], [0, 0, 3, ..."
4,"{'directed': True, 'multigraph': False, 'graph...",What man was a famous American author and also...,Q5686,Q7245,False,"{'edge_index': [[1, 3, 4, 4, 4], [0, 0, 0, 2, ..."


In [25]:
res_csv = pd.read_csv(
        f"/workspace/storage/mintaka_seq2seq/{dataset_type}/test/results.csv"
)

In [26]:
class EvalGraphDataset(Dataset):
    def __init__(self, is_corrects, graphs):
        self.data = []
        self.correct = []
        for is_correct, graph in zip(is_corrects, graphs):
            preproc_graph = preprocess(graph)
            if preproc_graph['input_edges'].shape[2] != 0:
                self.data.append(preproc_graph)
                self.correct.append(is_correct)
    
    def get_new_correct(self):
        return self.correct
    
    def __len__(self):
        return len(self.data)

    def __getitem__(self, idx):
        return self.data[idx]

In [27]:
final_acc, top200_total, top1_total, seq2seq_correct = 0, 0, 0, 0
    
for idx, group in tqdm(res_csv.iterrows()):
    curr_question_df = test_df[test_df["question"] == group['question']]
    if len(curr_question_df) == 0: # we don't have subgraph for this question, take answer from seq2seq
        if group["answer_0"] == group["target"]:
            seq2seq_correct += 1
        else: # check if answer exist in 200 beams for question with no subgraphs
            all_beams = group.tolist()[2:-1] # all 200 beams
            all_beams = list(set(all_beams))
            top200_total += 1 if group["target"] in all_beams else 0
            
    else: # we have subgraph for this question  
        all_beams = group.tolist()[2:-1] # all 200 beams
        all_beams = list(set(all_beams))
        
        if group["target"] not in all_beams: # no correct answer in beam
            continue
            
        # correct answer exist in beam
        top1_total += 1 if group["answer_0"] == group["target"] else 0
        top200_total += 1
        
        transformed_graphs = curr_question_df["transformed_graph"].tolist()
        is_corrects = curr_question_df["correct"].tolist()
        current_dataset = EvalGraphDataset(is_corrects, transformed_graphs)
        filtered_is_correct = current_dataset.get_new_correct()
        
        current_dataloader = torch.utils.data.DataLoader(current_dataset, 
                                                         batch_size=len(transformed_graphs), 
                                                         collate_fn=data_collator, 
                                                         shuffle=False)

        # batch size should only be one
        for item in current_dataloader:
            logits = outputs = model(input_nodes = item['input_nodes'].to(device), 
                                    input_edges = item['input_edges'].to(device),
                                    attn_bias = item['attn_bias'].to(device),
                                    in_degree = item['in_degree'].to(device),
                                    out_degree = item['out_degree'].to(device),
                                    spatial_pos = item['spatial_pos'].to(device),
                                    attn_edge_type = item['attn_edge_type'].to(device))
            mse_pred = outputs.logits.flatten()
            max_idx = mse_pred.argmax()
        
        if filtered_is_correct[max_idx]:
            final_acc += 1 
              

# final rerankinga, top1 and top200 result
reranking_res = (final_acc + seq2seq_correct)/ len(res_csv)
top200 = (top200_total + seq2seq_correct)/len(res_csv)
top1 = (top1_total + seq2seq_correct)/ len(res_csv)

        

4000it [00:25, 156.12it/s]


In [28]:
top1, reranking_res, top200

(0.25425, 0.22, 0.64375)