File size: 5,873 Bytes
73df34b
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
import os
import json
import argparse
from tqdm import tqdm

import torch
import torch.distributed as dist
from torch.utils.data import DataLoader

from data.dataset_for_clean_descrip import PoseHICODetDataset
from data.convsersation import Conversation_For_Clean_Descrption as Conversation

from dataclasses import dataclass

from tools.vlm_backend import build_batch_tensors, decode_generated_text, load_model_and_processor

class StreamingJsonArrayWriter:
    def __init__(self, output_path):
        self.output_path = output_path
        self.file = None
        self.is_first = True

    def __enter__(self):
        self.file = open(self.output_path, "w", encoding="utf-8")
        self.file.write("[\n")
        self.file.flush()
        return self

    def write(self, item):
        if not self.is_first:
            self.file.write(",\n")
        json.dump(item, self.file, ensure_ascii=False, indent=2)
        self.file.flush()
        self.is_first = False

    def __exit__(self, exc_type, exc_val, exc_tb):
        if self.file is not None:
            self.file.write("\n]\n")
            self.file.close()


def disable_torch_init():
    """
    Disable the redundant torch default initialization to accelerate model creation.
    """
    setattr(torch.nn.Linear, "reset_parameters", lambda self: None)
    setattr(torch.nn.LayerNorm, "reset_parameters", lambda self: None)


@dataclass
class DataCollatorForSupervisedDataset(object):
    def __init__(self, processor, data_path):
        self.processor = processor
        self.conv = Conversation(
            system='',
            data_path=data_path
        )

    def __call__(self, data_dicts):
        batch_prompts = []
        batch_images = []
        result_meta = []

        for data_dict in data_dicts:
            batch_images.append(data_dict['image'])
            batch_prompts.append(self.conv.get_prompt(data_dict['meta']))
            result_meta.append(data_dict['meta'])

        batch_tensors = build_batch_tensors(
            processor=self.processor,
            prompts=batch_prompts,
            images=batch_images,
            system_prompt=self.conv.system,
        )
        return batch_tensors, result_meta


@torch.no_grad()
def worker(model, processor, dataset, args):
    rank = int(os.environ["LOCAL_RANK"])
    world_size = int(os.environ["WORLD_SIZE"])
    indices = list(range(rank, len(dataset), world_size))
    print("==>" + " Worker {} Started, responsible for {} images".format(rank, len(indices)))

    sub_dataset = torch.utils.data.Subset(dataset, indices)
    data_loader = DataLoader(
        sub_dataset,
        batch_size=args.batch_size,
        shuffle=False,
        num_workers=0,
        collate_fn=DataCollatorForSupervisedDataset(processor, args.data_path)
    )
    output_path = os.path.join(args.output_dir, f'refine_labels_{rank}.json')

    with StreamingJsonArrayWriter(output_path) as writer:
        for batch_tensors, result_meta in tqdm(data_loader):
            input_ids = batch_tensors['input_ids'].cuda()
            batch_tensors = {k: v.cuda() for k, v in batch_tensors.items() if isinstance(v, torch.Tensor)}
            with torch.inference_mode():
                output_dict = model.generate(
                    do_sample=False,
                    output_scores=True,
                    return_dict_in_generate=True,
                    max_new_tokens=args.max_new_tokens,
                    output_logits=True,
                    **batch_tensors,
                )
                output_ids = output_dict['sequences']

            for input_id, output_id, meta in zip(input_ids, output_ids, result_meta):
                input_token_len = input_id.shape[0]
                n_diff_input_output = (input_id != output_id[:input_token_len]).sum().item()
                if n_diff_input_output > 0:
                    print(f'[Warning] Sample: {n_diff_input_output} output_ids are not the same as the input_ids')
                output = decode_generated_text(processor, output_id, input_id)
                meta['refined_description'] = output
                writer.write(meta)


def eval_model(args):
    dist.init_process_group(backend='nccl')
    rank = int(os.environ["LOCAL_RANK"])
    world_size = int(os.environ["WORLD_SIZE"])

    print('Init process group: world_size: {}, rank: {}'.format(world_size, rank))
    torch.cuda.set_device(rank)

    disable_torch_init()
    backend_name, model, processor = load_model_and_processor(
        model_path=args.model_path,
        backend=args.model_backend,
        torch_dtype=args.torch_dtype,
        trust_remote_code=True,
    )
    print(f'Using model backend: {backend_name}')
    model = model.cuda()
    model.eval()

    dataset = PoseHICODetDataset(
        data_path=args.data_path,
        multimodal_cfg=dict(
            image_folder=os.path.join(args.data_path, 'Images/images/train2015'),
            data_augmentation=False,
            image_size=336,
        ),
        annotation_path=args.annotation_path,
        max_samples=args.max_samples,
    )
    worker(model, processor, dataset, args)


if __name__ == "__main__":
    parser = argparse.ArgumentParser()
    parser.add_argument("--model-path", type=str, default="facebook/opt-350m")
    parser.add_argument("--data-path", type=str, default="")
    parser.add_argument("--annotation-path", type=str, default="./outputs/merged_labels.json")
    parser.add_argument("--output-dir", type=str, default="")
    parser.add_argument("--batch-size", type=int, default=8)
    parser.add_argument("--max-new-tokens", type=int, default=512)
    parser.add_argument("--max-samples", type=int, default=0)
    parser.add_argument("--model-backend", type=str, default="auto")
    parser.add_argument("--torch-dtype", type=str, default="bfloat16")
    args = parser.parse_args()

    eval_model(args)