reader.py 7.0 KB
Newer Older
S
sneaxiy 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17
# Copyright (c) 2019 PaddlePaddle Authors. All Rights Reserved.
#
# 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 core
import six
import threading
S
sneaxiy 已提交
18 19
from .framework import Program, Variable, program_guard, default_main_program, default_startup_program
from .executor import global_scope
S
sneaxiy 已提交
20
from .data_feeder import DataFeeder
S
sneaxiy 已提交
21 22
from .layers.io import monkey_patch_reader_methods, _copy_reader_var_, double_buffer
import unique_name
S
sneaxiy 已提交
23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41

__all__ = ['PyReader']


def _convert_places(places):
    if not isinstance(places, (list, tuple)):
        places = [places]

    ret = []
    for p in places:
        if not isinstance(p, core.Place):
            tmp = core.Place()
            tmp.set_place(p)
            p = tmp

        ret.append(p)
    return ret


S
sneaxiy 已提交
42 43 44 45 46 47 48 49
class PyReader(object):
    unique_name_generator = unique_name.UniqueNameGenerator()

    def __init__(self,
                 feed_list,
                 capacity,
                 use_double_buffer=True,
                 iterable=True):
S
sneaxiy 已提交
50 51
        self._tensor_reader = None
        self._thread = None
S
sneaxiy 已提交
52 53 54
        self._iterable = iterable
        self._use_double_buffer = use_double_buffer
        self._capacity = capacity
S
sneaxiy 已提交
55
        self._feed_list = feed_list
S
sneaxiy 已提交
56 57 58
        self._scope = global_scope()
        if not self._iterable:
            self._init_non_iterable()
S
sneaxiy 已提交
59

S
sneaxiy 已提交
60 61
    def _init_iterable(self, places):
        self._var_names = [v.name for v in self._feed_list]
S
sneaxiy 已提交
62
        self._places = _convert_places(places)
S
sneaxiy 已提交
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
        self._queue = core.init_lod_tensor_blocking_queue(core.Variable(),
                                                          self._capacity)
        self._reader = core.create_py_reader(
            self.queue, self._var_names, self._places, self._use_double_buffer)

    def _init_non_iterable(self):
        lod_levels = []
        dtypes = []
        shape_concat = []
        ranks = []
        shapes = []

        for feed_data in self._feed_list:
            dtypes.append(feed_data.dtype)
            shape_concat.extend(feed_data.shape)
            ranks.append(len(feed_data.shape))
            shapes.append(feed_data.shape)
            lod_levels.append(feed_data.lod_level)

        queue_name = PyReader.unique_name_generator('lod_tensor_blocking_queue')
        reader_name = PyReader.unique_name_generator('create_py_reader')
        double_buffer_name = PyReader.unique_name_generator('double_buffer')

        var = self._scope.var(queue_name)
        self._queue = core.init_lod_tensor_blocking_queue(var, self._capacity)

        startup_blk = default_startup_program().current_block()
        startup_var = startup_blk.create_var(name=reader_name)

        startup_blk.append_op(
            type='create_py_reader',
            inputs={'blocking_queue': [queue_name]},
            outputs={'Out': [startup_var]},
            attrs={
                'shape_concat': shape_concat,
                'lod_levels': lod_levels,
                'ranks': ranks
            })

        startup_var.desc.set_dtypes(dtypes)
        startup_var.persistable = True

        main_prog_var = _copy_reader_var_(
            default_main_program().current_block(), startup_var)

        main_prog_var.stop_gradient = True
        main_prog_var.persistable = True

        reader = monkey_patch_reader_methods(main_prog_var)
        if self._use_double_buffer:
            double_buffer_reader = double_buffer(
                reader, name=double_buffer_name)
            # we return a double buffer reader. However, the reset method comes from
            # py_reader.
            double_buffer_reader.reset = reader.reset
            reader = double_buffer_reader

        self._reader = reader

        default_main_program().current_block().append_op(
            type='read',
            inputs={'Reader': [self._reader]},
            outputs={'Out': self._feed_list})

    @property
    def queue(self):
        return self._queue

    @property
    def iterable(self):
        return self._iterable
S
sneaxiy 已提交
134 135

    def __call__(self):
S
sneaxiy 已提交
136
        assert self.iterable, "PyReader is not iterable"
S
sneaxiy 已提交
137 138 139 140 141
        assert self._tensor_reader is not None, \
            "Data source of PyReader has not set yet"

        class Iterator(object):
            def __init__(self, reader):
S
sneaxiy 已提交
142 143
                self._reader = reader._reader
                self._reset = reader._reset
S
sneaxiy 已提交
144 145 146 147 148

            def __iter__(self):
                return self

            def next(self):
S
sneaxiy 已提交
149
                ret = self._reader.read_next()
S
sneaxiy 已提交
150
                if ret:
S
sneaxiy 已提交
151 152
                    return ret
                else:
S
sneaxiy 已提交
153
                    self._reset()
S
sneaxiy 已提交
154 155
                    raise StopIteration

S
sneaxiy 已提交
156
        self._start()
S
sneaxiy 已提交
157 158
        return Iterator(self)

S
sneaxiy 已提交
159
    def _reset(self):
S
sneaxiy 已提交
160 161 162 163 164 165
        self._reader.reset()
        self._thread.join()

    def start(self):
        assert not self._iterable, "start() cannot be called when PyReader is iterable"
        self._start()
S
sneaxiy 已提交
166

S
sneaxiy 已提交
167 168 169 170 171
    def reset(self):
        assert not self._iterable, "reset() cannot be called when PyReader is iterable"
        self._reset()

    def _start(self):
S
sneaxiy 已提交
172 173 174 175 176 177 178 179 180 181 182
        def __thread_main__():
            for tensors in self._tensor_reader():
                array = core.LoDTensorArray()
                for item in tensors:
                    if not isinstance(item, core.LoDTensor):
                        tmp = core.LoDTensor()
                        tmp.set(item, core.CPUPlace())
                        item = tmp

                    array.append(item)

S
sneaxiy 已提交
183
                if not self._queue.push(array):
S
sneaxiy 已提交
184 185
                    break

S
sneaxiy 已提交
186
            self._queue.close()
S
sneaxiy 已提交
187 188 189 190 191

        self._thread = threading.Thread(target=__thread_main__)
        self._thread.daemon = True
        self._thread.start()

S
sneaxiy 已提交
192
    def decorate_paddle_reader(self, reader, places=None):
S
sneaxiy 已提交
193 194 195 196 197 198 199 200 201 202 203
        assert self._tensor_reader is None, \
            "Cannot reset the data source of PyReader"
        with program_guard(Program(), Program()):
            feeder = DataFeeder(
                feed_list=self._feed_list, place=core.CPUPlace())
            paddle_reader = feeder.decorate_reader(reader, multi_devices=False)

        def __tensor_reader_impl__():
            for slots in paddle_reader():
                yield [slots[var.name] for var in self._feed_list]

S
sneaxiy 已提交
204
        self.decorate_tensor_provider(__tensor_reader_impl__, places)
S
sneaxiy 已提交
205

S
sneaxiy 已提交
206
    def decorate_tensor_provider(self, reader, places=None):
S
sneaxiy 已提交
207 208 209
        assert self._tensor_reader is None, \
            "Cannot reset the data source of PyReader"
        self._tensor_reader = reader
S
sneaxiy 已提交
210 211 212
        if self._iterable:
            assert places is not None, "Places cannot be None when py_reader is iterable"
            self._init_iterable(places)