import logging import time from concurrent.futures import ThreadPoolExecutor from dataclasses import dataclass from app.connectors.database.snowflake_client import SnowflakeClient from app.connectors.database.types import EmailDomainValidation from app.services import EmailDomainValidationService logger = logging.getLogger(__name__) @dataclass(kw_only=True) class ValidateEmailDomainsHandler: email_domain_validation_service: EmailDomainValidationService snowflake_client: SnowflakeClient domains_chunk_size: int threads_count: int def handle(self) -> None: t0 = time.time() logger.info("Domains processing started") domains = self.snowflake_client.get_domains_to_verify(self.domains_chunk_size) with ThreadPoolExecutor(max_workers=self.threads_count) as pool: for domains_chunk in domains: results = list(pool.map(self.validate_domain, domains_chunk)) self.snowflake_client.update_domains(results) logger.info("Domains processing finished") logger.debug("Spent %s seconds", int(time.time() - t0)) def validate_domain(self, domain: EmailDomainValidation) -> EmailDomainValidation: try: logger.info("Processing %s domain", domain.email_domain) self.email_domain_validation_service.validate_email_domain( domain.email_domain ) domain.is_valid = True domain.attempts = 0 domain.reason = None except Exception as exc: logger.warning( '%s domain validation fail. Reason: "%s"', domain.email_domain, exc ) domain.is_valid = False domain.attempts += 1 domain.reason = str(exc) return domain