【问题标题】:Import a JSON project wise, so it loads just once明智地导入 JSON 项目,因此它只加载一次
【发布时间】:2020-11-25 08:54:13
【问题描述】:

我有一个 Python 项目,它针对特定模式执行 JSON 验证。 它将作为 GCP Dataflow 中的转换步骤运行,因此在运行之前收集所有依赖项非常重要,以避免一次又一次地下载相同的文件。

架构放置在单独的 Git 存储库中。 Transformer 的本质是您在课堂上收到一条记录,然后您就可以使用它。典型的流程是加载 JSON 模式,根据它验证记录,然后对无效和有效的内容进行处理。以这种方式加载模式意味着我从存储库中下载每条记录的模式,它可能是数十万。 代码被“克隆”到工作人员中,然后独立工作。

受到 Python 在开始时(一次)加载需求并将它们用作导入的方式的启发,我想我可以将存储库(JSON 模式所在的位置)添加为 Python 需求,然后简单地在我的 Python 代码。但当然,它是一个 JSON,而不是要导入的 Python 模块。它是如何工作的?

一个例子是这样的:

  • requirements.txt
git+git://github.com/path/to/json/schema@41b95ec
  • dataflow_transformer.py
import apache_beam as beam
import the_downloaded_schema
from jsonschema import validate

class Verifier(beam.DoFn):

    def process(self, record: dict):
        validate(instance=record, schema=the_downloaded_schema)

        # ... more stuff

        yield record

class Transformer(beam.PTransform):
    def expand(self, record):
        return (
            record
            | "Verify Schema" >> beam.ParDo(Verifier())
        )

【问题讨论】:

  • 您是否考虑过使用侧输入?您可以将 JSON 模式作为辅助输入传递,并使用它在主输入中验证您的记录。 Here 是它的文档。对你有帮助吗?
  • @AlexandreMoraes 我去看看,听起来很有趣,谢谢!

标签: python json google-cloud-dataflow apache-beam jsonschema


【解决方案1】:

您可以加载一次 json 架构并将其用作辅助输入。 一个例子:

import json
import requests

json_current='https://covidtracking.com/api/v1/states/current.json'

def get_json_schema(url):
  with requests.Session() as session:
    schema = json.loads(session.get(url).text)
  return schema

schema_json = get_json_schema(json_current)

def feed_schema(data, schema):
  yield {'record': data, 'schema': schema[0]}

schema = p | beam.Create([schema_json])
data = p | beam.Create(range(10))
data_with_schema = data | beam.FlatMap(feed_schema, schema=beam.pvalue.AsSingleton(schema))
# Now do your schema validation

只是演示data_with_schema pcollection 的样子

【讨论】:

  • 我最终实现了侧输入,以此为灵感。按预期工作。谢谢你:)
【解决方案2】:

为什么不直接使用一个类来加载使用缓存的资源以防止重复加载?大致如下:

class JsonLoader:
    def __init__(self):
        self.cache = set()

    def import(self, filename):
        filename = os.path.absname(filename)
        if filename not in self.cache:
            self._load_json(filename)
            self.cache.add(filename)

    def _load_json(self, filename):
        ...

【讨论】:

  • 这样,它会为每个工作人员下载一次架构。并且会在流程结束时被销毁,因此架构上的任何更新都将在管道的每次新运行中加载。很酷,谢谢 :-)
  • 我不太明白:你说你想使用某事。类似于 Python 的 import 机制,这表明您正在运行单个 python 进程。但是你也可以只使用上面那个类的单例实例?!
  • 我想要实现的是避免从每条记录中从 git 下载模式。我的第一印象是尽可能使用类似于 Python 的 import 的东西。缓存调用的单例会大大减少下载,然后我不会丢弃,但仍打算采用我的第一个方法。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-04-14
  • 2016-03-27
  • 1970-01-01
  • 2016-10-02
相关资源
最近更新 更多