【问题标题】:How to use Google App Engine MapReduce with large amount of data ( >1000000 )如何使用 Google App Engine MapReduce 处理大量数据(>1000000)
【发布时间】:2017-01-27 07:05:53
【问题描述】:

我实现了非常简单的 mapreduce 管道,但遇到了一些问题。

场景是,云数据存储中的一种模型中有超过 1000000 个实体,我想检查所有实体是否每个实体都具有不一致的属性。

这是我的代码 sn-p。

class User(ndb.model)
    parent = ndb.KeyProperty(Group) # want to check if this key property actually exist


class CheckKeyExistencePipeline(pipeline.Pipeline):

    def map(self, entity):
        logging.info(entity.urlsafe()) # added for debug
        prop = getattr(entity, 'parent')
        if not prop.get():
            yield 'parent does not exist: %s\n' % (entity.key.urlsafe())

    def run(self, modelname, shards):
        mapreduce_pipeline.MapperPipeline(
            'parent check',
            handler_spec='CheckKeyExistencePipeline.map',
            input_reader_spec='mapreduce.input_readers.DatastoreInputReader',
            output_writer_spec="mapreduce.output_writers.GoogleCloudStorageOutputWriter",
            params={
                'input_reader': {
                    'entity_kind': 'User',
                },
                'output_writer': {
                    'bucket_name': app_identity.get_default_gcs_bucket_name(),
                    'content_type': 'text/plain'
                }
            },
            shards=10)

问题是,它确实经常显示如下错误。

超过 128 MB 的软专用内存限制,之后为 133 MB 总共处理 2 个请求

当我使用大约 10000 个实体的数据运行此代码时,没有问题。 有什么问题?如何正确配置此管道以应用大量数据?

EDIT1

我修改为不使用ndb缓存,但似乎没有任何改善。根据源代码,我猜缓存已经默认关闭了。

https://github.com/GoogleCloudPlatform/appengine-mapreduce/blob/6e103ac52855c3214de3ec3721d6ec0e7edd5f77/python/src/mapreduce/util.py#L381-L383

def _set_ndb_cache_policy():
    """Tell NDB to never cache anything in memcache or in-process.
    This ensures that entities fetched from Datastore input_readers via NDB
    will not bloat up the request memory size and Datastore Puts will avoid
    doing calls to memcache. Without this you get soft memory limit exits,
    which hurts overall throughput.
    """
    ndb_ctx = ndb.get_context()
    ndb_ctx.set_cache_policy(lambda key: False)
    ndb_ctx.set_memcache_policy(lambda key: False)

我做了进一步的调查以找出问题所在。我将其中一个映射器参数processing_rate 设置为 10,shards 设置为 100,以便它只为每个任务处理 1 或 2 个实体。

这是 mapreduce 统计数据。该图似乎合理。 (此时管道还没有完成。)

但是当我检查其中一个工作任务的跟踪日志时,真的很奇怪。它显示了一堆/datastore_v3.Next/datastore_v3.Get 尽管事实上“map”函数只被调用了两次(根据我的调试日志。)因为我没有改变batch_size,它应该是50。所以,在我的理解, /datastore_v3.Next 应该只被调用一次,/datastore_v3.Get 应该被调用两次。

有人知道为什么会触发这么多对数据库的 RPC 调用吗?

EDIT2

我再次进行了进一步调查并简化了代码。 map 函数只是通过调用get 函数来使用ndb.Key 获取数据。

class CheckKeyExistencePipeline(pipeline.Pipeline):

    def map(self, entity):
        logging.info('start')
        entity.parent.get()
        logging.info('end')


    def run(self):
        mapreduce_pipeline.MapperPipeline(
            'parent check',
            handler_spec='CheckKeyExistencePipeline.map',
            input_reader_spec='mapreduce.input_readers.DatastoreInputReader',
            output_writer_spec="mapreduce.output_writers.GoogleCloudStorageOutputWriter",
            params={
                'input_reader': {
                    'entity_kind': 'User',
                },
                'output_writer': {
                    'bucket_name': app_identity.get_default_gcs_bucket_name(),
                    'content_type': 'text/plain'
                }
            },
            shards=10)

Stackdriver的Tracelog是这样的。

它只是调用get,但它在“开始”和“结束”之间多次触发 RPC 调用。这似乎有点奇怪,可能是这种内存消耗的原因之一。这是正常行为吗?

【问题讨论】:

  • 有时您可以使用ndb 提供的每个实例缓存来填充大量内存。您可以在不使用缓存的情况下获取实体:prop.get(use_cache=False)。在这种情况下,由于您没有连续对同一个键执行一堆获取操作,因此无论如何您都不会从缓存中获得任何东西......您也可以尝试increase the instance class 以获取处理这些请求的服务。 ..
  • 是的,缓存最终将成为这里的关键组件(通常缓存 ndb 实体没有帮助,除非您知道您将重用它们)。您可能希望将迭代器改为基于键的,以便更好地控制缓存行为。
  • 哦,我明白了。我没有考虑缓存。我会尝试。谢谢!
  • 杰夫,“基于密钥”是什么意思?

标签: python google-app-engine mapreduce


【解决方案1】:

我想我终于找到了问题所在。

关键是 MapReduce 使用 ndb.query.iter 并使用 eventloops 来管理异步 RPC 调用。在我的 MapReduce 调用中,它触发了两种 RPC 调用,一种是由 MapReduce 库触发以获取数据库记录(A),另一种是由我的map 函数(B)触发。

如果我没有在 map 函数中触发任何 RPC 调用,则没有地方可以触发下一个 RPC 调用。这意味着下一个 (A) 仅在批量迭代 50 条记录后触发。但是,(B) 的尝试触发了下一个 RPC 调用,并且由于 RPC 调用不是串行触发的(这意味着它不是 FIFO 队列),因此很可能会连续触发 (A) 直到它获取所有实体.

我将shards 设置为 100,但仍然有一个分片负责总共 10000 条记录。因此,这超出了软内存限制。

当我将shards 增加到 10000 时,会发生另一个错误...

总之,没有办法使用 MapReduce 处理大数据和低内存实例。我猜。

详情请查看以下问题。

https://code.google.com/p/googleappengine/issues/detail?id=11648 https://code.google.com/p/googleappengine/issues/detail?id=9610(原创)

【讨论】:

    猜你喜欢
    • 2011-04-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-06-22
    • 2020-02-10
    相关资源
    最近更新 更多