From 8084bc9e9c297122a5cb9b2a1928b498fa2f8288 Mon Sep 17 00:00:00 2001 From: sondrfos Date: Tue, 8 Sep 2026 11:57:49 +0000 Subject: [PATCH 1/2] add named loggers --- datacontract/api.py | 14 +++++----- datacontract/catalog/catalog.py | 4 ++- datacontract/data_contract.py | 4 ++- .../fastjsonschema/check_jsonschema.py | 10 ++++--- .../fastjsonschema/s3/s3_read_files.py | 4 ++- .../engines/ibis/connections/kafka.py | 4 ++- datacontract/export/html_exporter.py | 4 ++- datacontract/imports/powerbi_importer.py | 4 ++- datacontract/imports/sql_importer.py | 6 +++-- datacontract/init/init_template.py | 4 ++- datacontract/lint/resolve.py | 26 ++++++++++--------- datacontract/lint/schema.py | 8 +++--- datacontract/model/run.py | 8 +++--- 13 files changed, 63 insertions(+), 37 deletions(-) diff --git a/datacontract/api.py b/datacontract/api.py index 80208003d..c796c66e3 100644 --- a/datacontract/api.py +++ b/datacontract/api.py @@ -16,6 +16,8 @@ from datacontract.model.exceptions import DataContractException from datacontract.model.run import Run +logger = logging.getLogger(__name__) + DATA_CONTRACT_EXAMPLE_PAYLOAD = """apiVersion: v3.1.0 kind: DataContract id: orders @@ -254,21 +256,21 @@ def check_api_key(api_key_header: str | None): correct_api_key = os.getenv("DATACONTRACT_CLI_API_KEY") if correct_api_key is None or correct_api_key == "": - logging.info("Environment variable DATACONTRACT_CLI_API_KEY is not set. Skip API key check.") + logger.info("Environment variable DATACONTRACT_CLI_API_KEY is not set. Skip API key check.") return if api_key_header is None or api_key_header == "": - logging.info("The API key is missing.") + logger.info("The API key is missing.") raise HTTPException( status_code=status.HTTP_401_UNAUTHORIZED, detail="Missing API key. Use Header 'x-api-key' to provide the API key.", ) if api_key_header != correct_api_key: - logging.info("The provided API key is not correct.") + logger.info("The provided API key is not correct.") raise HTTPException( status_code=status.HTTP_403_FORBIDDEN, detail="The provided API key is not correct.", ) - logging.info("Request authenticated with API key.") + logger.info("Request authenticated with API key.") pass @@ -331,8 +333,8 @@ async def test( ] = None, ) -> Run: check_api_key(api_key) - logging.info("Testing data contract...") - logging.info(body) + logger.info("Testing data contract...") + logger.info(body) return DataContract( data_contract_str=body, server=server, publish_url=publish_url, fastapi_url=str(request.url) ).test() diff --git a/datacontract/catalog/catalog.py b/datacontract/catalog/catalog.py index 4d5b063e1..a182d039e 100644 --- a/datacontract/catalog/catalog.py +++ b/datacontract/catalog/catalog.py @@ -11,6 +11,8 @@ from datacontract.data_contract import DataContract from datacontract.export.html_exporter import get_version +logger = logging.getLogger(__name__) + def _get_owner(odcs: OpenDataContractStandard) -> Optional[str]: """Get the owner from ODCS customProperties or team.""" @@ -24,7 +26,7 @@ def _get_owner(odcs: OpenDataContractStandard) -> Optional[str]: def create_data_contract_html(contracts, file: Path, path: Path, schema: str): - logging.debug(f"Creating data contract html for file {file} and schema {schema}") + logger.debug(f"Creating data contract html for file {file} and schema {schema}") data_contract = DataContract( data_contract_file=f"{file.absolute()}", inline_references=True, schema_location=schema ) diff --git a/datacontract/data_contract.py b/datacontract/data_contract.py index e20faa6ba..01b43b272 100644 --- a/datacontract/data_contract.py +++ b/datacontract/data_contract.py @@ -19,6 +19,8 @@ from datacontract.model.exceptions import DataContractException, DataContractValidationErrors from datacontract.model.run import Check, ResultEnum, Run +logger = logging.getLogger(__name__) + class DataContract: def __init__( @@ -171,7 +173,7 @@ def test(self) -> Run: engine="datacontract", ) ) - logging.exception("Exception occurred") + logger.exception("Exception occurred") run.log_error(str(e)) run.finish() diff --git a/datacontract/engines/fastjsonschema/check_jsonschema.py b/datacontract/engines/fastjsonschema/check_jsonschema.py index 2d6dbe011..7aaffc7c3 100644 --- a/datacontract/engines/fastjsonschema/check_jsonschema.py +++ b/datacontract/engines/fastjsonschema/check_jsonschema.py @@ -14,6 +14,8 @@ from datacontract.model.exceptions import DataContractException from datacontract.model.run import Check, ResultEnum, Run +logger = logging.getLogger(__name__) + # Thread-safe cache for primaryKey fields. _primary_key_cache = {} _cache_lock = threading.Lock() @@ -88,13 +90,13 @@ def process_exceptions(run, exceptions: List[DataContractException]): def validate_json_stream( schema: dict, model_name: str, validate: Callable, json_stream: Generator[Any, Any, None] ) -> List[DataContractException]: - logging.info(f"Validating JSON stream for model: '{model_name}'.") + logger.info(f"Validating JSON stream for model: '{model_name}'.") exceptions: List[DataContractException] = [] for json_obj in json_stream: try: validate(json_obj) except JsonSchemaValueException as e: - logging.warning(f"Validation failed for JSON object with type: '{model_name}'.") + logger.warning(f"Validation failed for JSON object with type: '{model_name}'.") primary_key_value = get_primary_key_value(schema, model_name, json_obj) exceptions.append( DataContractException( @@ -108,7 +110,7 @@ def validate_json_stream( ) ) if not exceptions: - logging.info(f"All JSON objects in the stream passed validation for model: '{model_name}'.") + logger.info(f"All JSON objects in the stream passed validation for model: '{model_name}'.") return exceptions @@ -195,7 +197,7 @@ def process_local_file(run, server, schema, model_name, validate): ) for file in all_files: - logging.info(f"Processing file: {file}") + logger.info(f"Processing file: {file}") with open(file, "r") as f: process_json_file(run, schema, model_name, validate, f, server.delimiter) diff --git a/datacontract/engines/fastjsonschema/s3/s3_read_files.py b/datacontract/engines/fastjsonschema/s3/s3_read_files.py index 87447f2ed..2fc17dec6 100644 --- a/datacontract/engines/fastjsonschema/s3/s3_read_files.py +++ b/datacontract/engines/fastjsonschema/s3/s3_read_files.py @@ -4,13 +4,15 @@ from datacontract.model.exceptions import DataContractException from datacontract.model.run import ResultEnum +logger = logging.getLogger(__name__) + def yield_s3_files(s3_endpoint_url, s3_location): fs = s3_fs(s3_endpoint_url) files = fs.glob(s3_location) for file in files: with fs.open(file) as f: - logging.info(f"Downloading file {file}") + logger.info(f"Downloading file {file}") yield f.read() diff --git a/datacontract/engines/ibis/connections/kafka.py b/datacontract/engines/ibis/connections/kafka.py index 2b2234f41..207240604 100644 --- a/datacontract/engines/ibis/connections/kafka.py +++ b/datacontract/engines/ibis/connections/kafka.py @@ -10,6 +10,8 @@ from datacontract.model.exceptions import DataContractException from datacontract.model.run import ResultEnum +logger = logging.getLogger(__name__) + def _scala_binary_version() -> str: """Return the Scala binary version the installed PySpark was built against. @@ -108,7 +110,7 @@ def read_kafka_topic(spark, data_contract: OpenDataContractStandard, server: Ser model_name = schema_obj.name topic = schema_obj.physicalName or schema_obj.name - logging.info("Reading data from Kafka server %s topic %s", server.host, topic) + logger.info("Reading data from Kafka server %s topic %s", server.host, topic) df = ( spark.read.format("kafka") .options(**get_auth_options()) diff --git a/datacontract/export/html_exporter.py b/datacontract/export/html_exporter.py index c8fab8214..96629da33 100644 --- a/datacontract/export/html_exporter.py +++ b/datacontract/export/html_exporter.py @@ -10,6 +10,8 @@ from datacontract.export.exporter import Exporter from datacontract.export.mermaid_exporter import to_mermaid +logger = logging.getLogger(__name__) + class HtmlExporter(Exporter): def export(self, data_contract, schema_name, server, sql_server_type, export_args) -> str: @@ -64,5 +66,5 @@ def get_version() -> str: try: return version("datacontract_cli") except Exception as e: - logging.debug("Ignoring exception", e) + logger.debug("Ignoring exception", e) return "" diff --git a/datacontract/imports/powerbi_importer.py b/datacontract/imports/powerbi_importer.py index f7cdde8f2..d6392de88 100644 --- a/datacontract/imports/powerbi_importer.py +++ b/datacontract/imports/powerbi_importer.py @@ -19,6 +19,8 @@ from datacontract.imports.odcs_helper import create_odcs, create_property, create_schema_object, create_server from datacontract.model.exceptions import DataContractException +logger = logging.getLogger(__name__) + # --------------------------------------------------------------------------- # Power BI data type → (ODCS logical type, optional format) # --------------------------------------------------------------------------- @@ -217,7 +219,7 @@ def _build_odcs(bim: dict[str, Any], model_name: str) -> OpenDataContractStandar _apply_bim_relationships(bim_relationships, table_name_to_obj) if not schema_objects: - logging.warning("Power BI import produced an empty contract: No tables were found in the semantic model.") + logger.warning("Power BI import produced an empty contract: No tables were found in the semantic model.") schema_objects.sort(key=lambda s: s.name.lower()) odcs.schema_ = schema_objects diff --git a/datacontract/imports/sql_importer.py b/datacontract/imports/sql_importer.py index 4cb5c72cf..cf742affe 100644 --- a/datacontract/imports/sql_importer.py +++ b/datacontract/imports/sql_importer.py @@ -17,6 +17,8 @@ from datacontract.model.exceptions import DataContractException from datacontract.model.run import ResultEnum +logger = logging.getLogger(__name__) + class SqlDialect(str, Enum): postgres = "postgres" @@ -44,7 +46,7 @@ def import_sql(source: str, import_args: dict = None) -> OpenDataContractStandar try: parsed = sqlglot.parse_one(sql=sql, read=dialect) except Exception as e: - logging.error(f"Error sqlglot SQL: {str(e)}") + logger.error(f"Error sqlglot SQL: {str(e)}") raise DataContractException( type="import", name=f"Reading source from {source}", @@ -60,7 +62,7 @@ def import_sql(source: str, import_args: dict = None) -> OpenDataContractStandar if server_type is not None: server_defaults = get_server_defaults(server_type) odcs.servers = [create_server(name=server_type, server_type=server_type, **server_defaults)] - logging.warning( + logger.warning( "SQL import generated a server block with placeholder connection values. " "Update host, port, database, and schema in the output before use." ) diff --git a/datacontract/init/init_template.py b/datacontract/init/init_template.py index 7578e51b5..9dad1a0e8 100644 --- a/datacontract/init/init_template.py +++ b/datacontract/init/init_template.py @@ -5,10 +5,12 @@ DEFAULT_DATA_CONTRACT_INIT_TEMPLATE = "odcs-3.1.0.init.yaml" +logger = logging.getLogger(__name__) + def get_init_template(location: str = None) -> str: if location is None: - logging.info("Use default bundled template " + DEFAULT_DATA_CONTRACT_INIT_TEMPLATE) + logger.info("Use default bundled template " + DEFAULT_DATA_CONTRACT_INIT_TEMPLATE) schemas = resources.files("datacontract") template = schemas.joinpath("schemas", DEFAULT_DATA_CONTRACT_INIT_TEMPLATE) with template.open("r") as file: diff --git a/datacontract/lint/resolve.py b/datacontract/lint/resolve.py index 2211440a4..39d09caf0 100644 --- a/datacontract/lint/resolve.py +++ b/datacontract/lint/resolve.py @@ -18,6 +18,8 @@ from datacontract.model.odcs import is_open_data_contract_standard, is_open_data_product_standard from datacontract.model.run import ResultEnum +logger = logging.getLogger(__name__) + class _LaxOpenDataContractStandard(OpenDataContractStandard): """ODCS variant that accepts unknown top-level fields. @@ -66,7 +68,7 @@ def _resolve_jsonschema_compliance_error_message_path(yaml_str, message): f"properties.{yaml_str['schema'][int(schema_index)]['properties'][int(property_index)]['name']}", ) except Exception: - logging.warning("YAML doesn't conform to JSON schema. Could not resolve indexed schema or property names.") + logger.warning("YAML doesn't conform to JSON schema. Could not resolve indexed schema or property names.") except_message = message finally: return except_message @@ -295,7 +297,7 @@ def _definition_resolution_error( url: str, target_url: str, detail: str, original_exception: Exception | None = None ) -> DataContractException: reason = f"Could not resolve business definition '{url}' from {target_url}: {detail}" - logging.warning(reason) + logger.warning(reason) return DataContractException( type="lint", result=ResultEnum.failed, @@ -312,7 +314,7 @@ def _resolve_data_contract_from_str( yaml_dict = _to_yaml(data_contract_str) if is_open_data_product_standard(yaml_dict): - logging.info("Cannot import ODPS, as not supported") + logger.info("Cannot import ODPS, as not supported") raise DataContractException( type="schema", result=ResultEnum.failed, @@ -322,7 +324,7 @@ def _resolve_data_contract_from_str( ) if is_open_data_contract_standard(yaml_dict): - logging.info("Importing ODCS v3") + logger.info("Importing ODCS v3") # When a custom JSON schema is provided, treat it as the source of # truth and accept extra top-level fields the standard ODCS Pydantic # class would reject. @@ -337,7 +339,7 @@ def _resolve_data_contract_from_str( return odcs # For DCS format, we need to convert it to ODCS - logging.info("Importing DCS format - converting to ODCS") + logger.info("Importing DCS format - converting to ODCS") from datacontract.imports.dcs_importer import convert_dcs_to_odcs, parse_dcs_from_dict dcs = parse_dcs_from_dict(yaml_dict) @@ -366,7 +368,7 @@ def _to_yaml(data_contract_str) -> dict: try: return yaml.load(data_contract_str, Loader=_SafeLoaderNoTimestamp) except Exception as e: - logging.warning(f"Cannot parse YAML. Error: {str(e)}") + logger.warning(f"Cannot parse YAML. Error: {str(e)}") raise DataContractException( type="lint", result="failed", @@ -388,7 +390,7 @@ def _validation_error_to_exception(error_message: str, original_exception=None) def _validate_json_schema(yaml_str, schema_location: str | Path = None, all_errors: bool = False): - logging.debug(f"Linting data contract with schema at {schema_location}") + logger.debug(f"Linting data contract with schema at {schema_location}") schema = fetch_schema(schema_location) if all_errors: validator_cls = validators.validator_for(schema) @@ -396,20 +398,20 @@ def _validate_json_schema(yaml_str, schema_location: str | Path = None, all_erro validator = validator_cls(schema=schema) errors = sorted(validator.iter_errors(yaml_str), key=lambda error: list(error.path)) if errors: - logging.warning(f"Data Contract YAML is invalid. Validation errors: {len(errors)}") + logger.warning(f"Data Contract YAML is invalid. Validation errors: {len(errors)}") raise DataContractValidationErrors( [_validation_error_to_exception(error.message, original_exception=error) for error in errors] ) - logging.debug("YAML data is valid.") + logger.debug("YAML data is valid.") return try: fastjsonschema.validate(schema, yaml_str, use_default=False) - logging.debug("YAML data is valid.") + logger.debug("YAML data is valid.") except JsonSchemaValueException as e: except_message = _resolve_jsonschema_compliance_error_message_path(yaml_str, e.message) - logging.warning(f"Data Contract YAML is invalid. Validation error: {except_message}") + logger.warning(f"Data Contract YAML is invalid. Validation error: {except_message}") raise _validation_error_to_exception(except_message, original_exception=e) except Exception as e: - logging.warning(f"Data Contract YAML is invalid. Validation error: {str(e)}") + logger.warning(f"Data Contract YAML is invalid. Validation error: {str(e)}") raise _validation_error_to_exception(str(e), original_exception=e) diff --git a/datacontract/lint/schema.py b/datacontract/lint/schema.py index b0e7867aa..c37714b4a 100644 --- a/datacontract/lint/schema.py +++ b/datacontract/lint/schema.py @@ -12,6 +12,8 @@ DEFAULT_DATA_CONTRACT_SCHEMA = "datacontract-1.2.1.schema.json" +logger = logging.getLogger(__name__) + def fetch_schema(location: str | Path = None) -> Dict[str, Any]: """ @@ -33,7 +35,7 @@ def fetch_schema(location: str | Path = None) -> Dict[str, Any]: """ if location is None: - logging.info("Use default bundled schema " + DEFAULT_DATA_CONTRACT_SCHEMA) + logger.info("Use default bundled schema " + DEFAULT_DATA_CONTRACT_SCHEMA) schemas = resources.files("datacontract") schema_file = schemas.joinpath("schemas", DEFAULT_DATA_CONTRACT_SCHEMA) with schema_file.open("r") as file: @@ -43,7 +45,7 @@ def fetch_schema(location: str | Path = None) -> Dict[str, Any]: location_str = str(location) if location_str.startswith("http://") or location_str.startswith("https://"): - logging.debug(f"Downloading schema from {location_str}") + logger.debug(f"Downloading schema from {location_str}") response = requests.get(location_str) schema = response.json() else: @@ -56,7 +58,7 @@ def fetch_schema(location: str | Path = None) -> Dict[str, Any]: result=ResultEnum.error, ) - logging.debug(f"Loading JSON schema locally at {location}") + logger.debug(f"Loading JSON schema locally at {location}") with open(location, "r") as file: schema = json.load(file) diff --git a/datacontract/model/run.py b/datacontract/model/run.py index c4b014eb8..2abc2d96e 100644 --- a/datacontract/model/run.py +++ b/datacontract/model/run.py @@ -6,6 +6,8 @@ from pydantic import BaseModel +logger = logging.getLogger(__name__) + class ResultEnum(str, Enum): passed = "passed" @@ -79,15 +81,15 @@ def calculate_result(self): self.result = ResultEnum.unknown def log_info(self, message: str): - logging.info(message) + logger.info(message) self.logs.append(Log(level="INFO", message=message, timestamp=datetime.now(timezone.utc))) def log_warn(self, message: str): - logging.warning(message) + logger.warning(message) self.logs.append(Log(level="WARN", message=message, timestamp=datetime.now(timezone.utc))) def log_error(self, message: str): - logging.error(message) + logger.error(message) self.logs.append(Log(level="ERROR", message=message, timestamp=datetime.now(timezone.utc))) def pretty(self): From 73c861bee5b1479d32dac042e7ca62c65edf7fe6 Mon Sep 17 00:00:00 2001 From: sondrfos Date: Tue, 8 Sep 2026 12:46:17 +0000 Subject: [PATCH 2/2] changelog and format --- CHANGELOG.md | 3 +++ datacontract/model/run.py | 1 + tests/test_import_powerbi.py | 1 - 3 files changed, 4 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index bf05b96da..e572a6a2b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Changed +- Internal logging uses named module loggers instead of the root logger + ## [1.1.3] - 2026-09-03 ### Added diff --git a/datacontract/model/run.py b/datacontract/model/run.py index b8a52db1c..5c54bb31b 100644 --- a/datacontract/model/run.py +++ b/datacontract/model/run.py @@ -41,6 +41,7 @@ def setter(self, value): return property(getter, setter, doc=message) + logger = logging.getLogger(__name__) diff --git a/tests/test_import_powerbi.py b/tests/test_import_powerbi.py index 4687783f6..d0ecc8c11 100644 --- a/tests/test_import_powerbi.py +++ b/tests/test_import_powerbi.py @@ -409,7 +409,6 @@ def test_import_bim_calculated_table_physical_type(): def test_import_pbit_from_zip(tmp_path): - result = import_powerbi_from_file(PBIT_FIXTURE) assert result is not None