Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- import logging
- import os
- import sys
- from dataclasses import dataclass
- from typing import List
- from uuid import uuid4
- import structlog
- from bson.objectid import ObjectId
- from dotenv import load_dotenv
- from pymongo import uri_parser
- from pymongo.collection import Collection
- from pymongo.mongo_client import MongoClient
- from etl import model
- from etl.anonymization import AccessorsAnonymizer, Tokenizer
- from etl.bigquery import Bigquery, CommandTableFacade
- from etl.processing import Processor
- from etl.util import fatal
- FALLBACK_MONGO_URI = "mongodb://localhost:27017/we"
- execution_id = str(uuid4())
- def main():
- cfg = get_config()
- if cfg.is_deployment:
- configure_logging()
- log = structlog.stdlib.get_logger().bind(execution_id=cfg.execution_id, environment=cfg.environment)
- log.info(f"Discovered configuration: {cfg}")
- cfg.deployment_config_sanity_check(log)
- # MongoCatalog, Tokenizer, AccessorsAnonymizer, Bigquery and CommandTableFacade are "library" classes. Notices that these classes are instantiated only once.
- mongo = MongoCatalog(cfg.mongo_rs_core_uri, cfg.mongo_rs_accessors_uri)
- tokenizer = Tokenizer(cfg.tokenizer_url, cfg.tokenizer_user, cfg.tokenizer_pw)
- filter_accessors_for_users = user_filter_factory(mongo.get_collection("user"))
- accessors_anonymizer = AccessorsAnonymizer(tokenizer, filter_accessors_for_users)
- bq = Bigquery(log, cfg.dataset_id, cfg.full_scan, cfg.execution_id, cfg.environment)
- def with_prefix(base_name: str):
- return f"{cfg.table_prefix}_{base_name}"
- cmd_table = CommandTableFacade(log, bq, with_prefix("exec_commands"))
- cmd_table.create()
- processors = [
- Processor(
- log,
- bq,
- cmd_table,
- mongo.get_collection("installation"),
- model.Table(
- with_prefix("exported_installations"),
- "",
- model.installation_table_definition(accessors_anonymizer),
- ),
- "installation",
- mongo_cursor_batch_size=300,
- ),
- Processor(
- log,
- bq,
- cmd_table,
- mongo.get_collection("branch"),
- model.Table(
- with_prefix("exported_branches"),
- "",
- model.branch_table_definition(),
- ),
- "branch",
- mongo_cursor_batch_size=3000,
- ),
- Processor(
- log,
- bq,
- cmd_table,
- mongo.get_collection("space"),
- model.Table(
- with_prefix("exported_spaces"),
- "",
- model.space_table_definition(accessors_anonymizer),
- ),
- "space",
- mongo_cursor_batch_size=300,
- ),
- Processor(
- log,
- bq,
- cmd_table,
- mongo.get_collection("user"),
- model.Table(
- with_prefix("exported_users"),
- "",
- model.user_table_definition(tokenizer),
- ),
- "user",
- mongo_cursor_batch_size=1000,
- ),
- ]
- for p in processors:
- p.setup()
- for p in processors:
- p.export()
- log.info("Done")
Advertisement
Add Comment
Please, Sign In to add comment