data_manager.py 13.0 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20
# Copyright (c) 2020 VisualDL Authors. All Rights Reserve.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
#     http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
# =======================================================================

import threading
import random
import collections

DEFAULT_PLUGIN_MAXSIZE = {
走神的阿圆's avatar
走神的阿圆 已提交
21
    "scalar": 1000,
22
    "image": 10,
23
    "histogram": 100,
走神的阿圆's avatar
走神的阿圆 已提交
24
    "embeddings": 50000000,
走神的阿圆's avatar
走神的阿圆 已提交
25
    "audio": 10,
走神的阿圆's avatar
走神的阿圆 已提交
26
    "pr_curve": 300,
P
Peter Pan 已提交
27
    "roc_curve": 300,
走神的阿圆's avatar
走神的阿圆 已提交
28 29
    "meta_data": 100,
    "text": 10
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
}


class Reservoir(object):
    """A map-to-arrays dict, with deterministic Reservoir Sampling.

    Store each reservoir bucket by key, and each bucket is a list sampling
    with reservoir algorithm.
    """

    def __init__(self, max_size, seed=0):
        """Creates a new reservoir.

        Args:
            max_size: The number of values to keep in the reservoir for each tag,
                if max_size is zero, all values will be kept in bucket.
            seed: The seed to initialize a random.Random().
            num_item_index: The index of data to add.

        Raises:
          ValueError: If max_size is not a nonnegative integer.
        """
        if max_size < 0 or max_size != round(max_size):
            raise ValueError("Max_size must be nonnegative integer.")
        self._max_size = max_size
        self._buckets = collections.defaultdict(
56 57
            lambda: _ReservoirBucket(max_size=self._max_size,
                                     random_instance=random.Random(seed))
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
        )
        self._mutex = threading.Lock()

    @property
    def keys(self):
        """Return all keys in self._buckets.

        Returns:
            All keys in reservoir buckets.
        :return:
        """
        with self._mutex:
            return list(self._buckets.keys())

    def _exist_in_keys(self, key):
        """Determine if key exists.

        Args:
            key: Key to determine if exists.

        Returns:
            True if key exists in buckets.keys, otherwise False.
        """
        return True if key in self._buckets.keys() else False

83 84
    def exist_in_keys(self, run, tag):
        """Determine if run_tag exists.
85 86 87 88

        For usage habits of VisualDL, actually call self._exist_in_keys()

        Args:
89
            run: Identity of one tablet.
90 91 92
            tag: Identity of one record in tablet.

        Returns:
93
            True if run_tag exists in buckets.keys, otherwise False.
94
        """
95
        key = run + "/" + tag
96 97 98 99 100 101 102 103
        return self._exist_in_keys(key)

    def _get_num_items_index(self, key):
        keys = self.keys
        if key not in keys:
            raise KeyError("Key %s not in buckets.keys()" % key)
        return self._buckets[key].num_items_index

104 105
    def get_num_items_index(self, run, tag):
        key = run + "/" + tag
106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122
        return self._get_num_items_index(key)

    def _get_items(self, key):
        """Get items with tag "key"

        Args:
            key: Key to finding bucket in reservoir buckets.

        Returns:
            One bucket in reservoir buckets by key.
        """
        keys = self.keys
        with self._mutex:
            if key not in keys:
                raise KeyError("Key %s not in buckets.keys()" % key)
            return self._buckets[key].items

123 124
    def get_items(self, run, tag):
        """Get items with tag 'run_tag'
125

126 127 128
        For usage habits of VisualDL, actually call self._get_items()

        Args:
129
            run: Identity of one tablet.
130 131 132
            tag: Identity of one record in tablet.

        Returns:
133
            One bucket in reservoir buckets by run and tag.
134
        """
135
        key = run + "/" + tag
136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156
        return self._get_items(key)

    def _add_item(self, key, item):
        """Add a new item to reservoir buckets with given tag as key.

        If bucket with key has not yet reached full size, each item will be
        added.

        If bucket with key is full, each item will be added with same
        probability.

        Add new item to buckets will always valid because self._buckets is a
        collection.defaultdict.

        Args:
            key: Tag of one bucket to add new item.
            item: New item to add to bucket.
        """
        with self._mutex:
            self._buckets[key].add_item(item)

走神的阿圆's avatar
走神的阿圆 已提交
157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175
    def _add_scalar_item(self, key, item):
        """Add a new scalar item to reservoir buckets with given tag as key.

        If bucket with key has not yet reached full size, each item will be
        added.

        If bucket with key is full, each item will be added with same
        probability.

        Add new item to buckets will always valid because self._buckets is a
        collection.defaultdict.

        Args:
            key: Tag of one bucket to add new item.
            item: New item to add to bucket.
        """
        with self._mutex:
            self._buckets[key].add_scalar_item(item)

176
    def add_item(self, run, tag, item):
177 178 179 180 181
        """Add a new item to reservoir buckets with given tag as key.

        For usage habits of VisualDL, actually call self._add_items()

        Args:
182
            run: Identity of one tablet.
183 184 185
            tag: Identity of one record in tablet.
            item: New item to add to bucket.
        """
186
        key = run + "/" + tag
187 188
        self._add_item(key, item)

走神的阿圆's avatar
走神的阿圆 已提交
189 190 191 192 193 194 195 196 197 198 199 200 201
    def add_scalar_item(self, run, tag, item):
        """Add a new scalar item to reservoir buckets with given tag as key.

        For usage habits of VisualDL, actually call self._add_items()

        Args:
            run: Identity of one tablet.
            tag: Identity of one record in tablet.
            item: New item to add to bucket.
        """
        key = run + "/" + tag
        self._add_scalar_item(key, item)

202 203 204 205
    def _cut_tail(self, key):
        with self._mutex:
            self._buckets[key].cut_tail()

206
    def cut_tail(self, run, tag):
207
        """Pop the last item in reservoir buckets.
208

209 210 211 212
        Sometimes the tail of the retrieved data is abnormal 0. This
        method is used to handle this problem.

        Args:
213
            run: Identity of one tablet.
214 215
            tag: Identity of one record in tablet.
        """
216
        key = run + "/" + tag
217 218
        self._cut_tail(key)

219 220 221 222

class _ReservoirBucket(object):
    """Data manager for sampling data, use reservoir sampling.
    """
223

224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245
    def __init__(self, max_size, random_instance=None):
        """Create a _ReservoirBucket instance.

        Args:
            max_size: The maximum size of reservoir bucket. If max_size is
                zero, the bucket has unbounded size.
            random_instance: The random number generator. If not specified,
                default to random.Random(0)
            num_item_index: The index of data to add.

        Raises:
            ValueError: If args max_size is not a nonnegative integer.
        """
        if max_size < 0 or max_size != round(max_size):
            raise ValueError("Max_size must be nonnegative integer.")
        self._max_size = max_size
        self._random = random_instance if random_instance is not None else \
            random.Random(0)
        self._items = []
        self._mutex = threading.Lock()
        self._num_items_index = 0

走神的阿圆's avatar
走神的阿圆 已提交
246 247 248
        self.max_scalar = None
        self.min_scalar = None

249 250 251 252 253 254 255 256 257 258 259 260 261
    def add_item(self, item):
        """ Add an item to bucket, replacing an old item with probability.

        Use reservoir sampling to add a new item to sampling bucket,
        each item in a steam has same probability stay in the bucket.

        Args:
            item: The item to add to reservoir bucket.
        """
        with self._mutex:
            if len(self._items) < self._max_size or self._max_size == 0:
                self._items.append(item)
            else:
走神的阿圆's avatar
走神的阿圆 已提交
262
                r = self._random.randint(1, self._num_items_index)
263 264 265 266 267
                if r < self._max_size:
                    self._items.pop(r)
                    self._items.append(item)
            self._num_items_index += 1

走神的阿圆's avatar
走神的阿圆 已提交
268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300
    def add_scalar_item(self, item):
        """ Add an scalar item to bucket, replacing an old item with probability.

        Use reservoir sampling to add a new item to sampling bucket,
        each item in a steam has same probability stay in the bucket.

        Args:
            item: The item to add to reservoir bucket.
        """
        with self._mutex:
            if not self.max_scalar or self.max_scalar.value < item.value:
                self.max_scalar = item
            if not self.min_scalar or self.min_scalar.value > item.value:
                self.min_scalar = item

            if len(self._items) < self._max_size or self._max_size == 0:
                self._items.append(item)
            else:
                if item.id == self.min_scalar.id or item.id == self.max_scalar.id:
                    r = self._random.randint(1, self._max_size - 1)
                else:
                    r = self._random.randint(1, self._num_items_index)
                if r < self._max_size:
                    if self._items[r].id == self.min_scalar.id or self._items[r].id == self.max_scalar.id:
                        if r - 1 > 0:
                            r = r - 1
                        elif r + 1 < self._max_size:
                            r = r + 1
                    self._items.pop(r)
                    self._items.append(item)

            self._num_items_index += 1

301 302 303 304 305 306 307 308 309 310 311 312 313 314 315
    @property
    def items(self):
        """Get self._items

        Returns:
            All items.
        """
        with self._mutex:
            return self._items

    @property
    def num_items_index(self):
        with self._mutex:
            return self._num_items_index

316 317 318 319 320 321 322 323 324 325
    def cut_tail(self):
        """Pop the last item in reservoir buckets.

        Sometimes the tail of the retrieved data is abnormal 0. This
        method is used to handle this problem.
        """
        with self._mutex:
            self._items.pop()
            self._num_items_index -= 1

326 327 328 329

class DataManager(object):
    """Data manager for all plugin.
    """
330

331 332 333 334 335 336 337
    def __init__(self):
        """Create a data manager for all plugin.

        All kinds of plugin has own reservoir, stored in a dict with plugin
        name as key.

        """
338 339 340 341 342 343 344 345 346 347
        self._reservoirs = {
            "scalar":
            Reservoir(max_size=DEFAULT_PLUGIN_MAXSIZE["scalar"]),
            "histogram":
            Reservoir(max_size=DEFAULT_PLUGIN_MAXSIZE["histogram"]),
            "image":
            Reservoir(max_size=DEFAULT_PLUGIN_MAXSIZE["image"]),
            "embeddings":
            Reservoir(max_size=DEFAULT_PLUGIN_MAXSIZE["embeddings"]),
            "audio":
走神的阿圆's avatar
走神的阿圆 已提交
348 349
            Reservoir(max_size=DEFAULT_PLUGIN_MAXSIZE["audio"]),
            "pr_curve":
走神的阿圆's avatar
走神的阿圆 已提交
350
            Reservoir(max_size=DEFAULT_PLUGIN_MAXSIZE["pr_curve"]),
P
Peter Pan 已提交
351 352
            "roc_curve":
            Reservoir(max_size=DEFAULT_PLUGIN_MAXSIZE["roc_curve"]),
走神的阿圆's avatar
走神的阿圆 已提交
353
            "meta_data":
走神的阿圆's avatar
走神的阿圆 已提交
354 355 356
            Reservoir(max_size=DEFAULT_PLUGIN_MAXSIZE["meta_data"]),
            "text":
            Reservoir(max_size=DEFAULT_PLUGIN_MAXSIZE["text"])
357
        }
358 359
        self._mutex = threading.Lock()

360 361 362 363 364 365 366 367 368 369 370 371 372 373 374
    def add_reservoir(self, plugin):
        """Add reservoir to reservoirs.

        Every reservoir is attached to one plugin.

        Args:
            plugin: Key to get one reservoir bucket for one specified plugin.
        """
        with self._mutex:
            if plugin not in self._reservoirs.keys():
                self._reservoirs.update({
                    plugin:
                    Reservoir(max_size=DEFAULT_PLUGIN_MAXSIZE[plugin])
                })

375 376 377 378 379 380 381 382 383 384 385 386 387 388
    def get_reservoir(self, plugin):
        """Get reservoir by plugin as key.

        Args:
            plugin: Key to get one reservoir bucket for one specified plugin.

        Returns:
            Reservoir bucket for plugin.
        """
        with self._mutex:
            if plugin not in self._reservoirs.keys():
                raise KeyError("Key %s not in reservoirs." % plugin)
            return self._reservoirs[plugin]

389
    def add_item(self, plugin, run, tag, item):
390 391
        """Add item to one plugin reservoir bucket.

392
        Use 'run', 'tag' for usage habits of VisualDL.
393 394 395

        Args:
            plugin: Key to get one reservoir bucket.
396
            run: Each tablet has different 'run'.
397 398 399 400
            tag: Tag will be used to generate paths of tablets.
            item: The item to add to reservoir bucket.
        """
        with self._mutex:
走神的阿圆's avatar
走神的阿圆 已提交
401 402 403 404
            if 'scalar' == plugin:
                self._reservoirs[plugin].add_scalar_item(run, tag, item)
            else:
                self._reservoirs[plugin].add_item(run, tag, item)
405 406 407 408 409 410 411 412 413 414 415 416

    def get_keys(self):
        """Get all plugin buckets name.

        Returns:
            All plugin keys.
        """
        with self._mutex:
            return self._reservoirs.keys()


default_data_manager = DataManager()