Skip to content

OMOP

Loader

pluginlake.omop.loader

OMOP data loading functions.

load_omop_table(file_path, table_name, *, validate=None, encoding=None)

Load OMOP CSV into validated Polars DataFrame.

Parameters:

Name Type Description Default
file_path Path

Path to CSV file.

required
table_name str

OMOP table name (e.g., 'person', 'condition_occurrence').

required
validate bool | None

Run schema validation. Uses config default if None.

None
encoding str | None

CSV encoding. Uses config default if None.

None

Returns:

Type Description
DataFrame

Polars DataFrame with validated schema.

Raises:

Type Description
FileNotFoundError

If file doesn't exist.

ValueError

If validation fails and skip_invalid_rows is False.

Source code in src/pluginlake/omop/loader.py
def load_omop_table(
    file_path: Path,
    table_name: str,
    *,
    validate: bool | None = None,
    encoding: str | None = None,
) -> pl.DataFrame:
    """Load OMOP CSV into validated Polars DataFrame.

    Args:
        file_path: Path to CSV file.
        table_name: OMOP table name (e.g., 'person', 'condition_occurrence').
        validate: Run schema validation. Uses config default if None.
        encoding: CSV encoding. Uses config default if None.

    Returns:
        Polars DataFrame with validated schema.

    Raises:
        FileNotFoundError: If file doesn't exist.
        ValueError: If validation fails and skip_invalid_rows is False.
    """
    settings = get_omop_settings()
    validate = validate if validate is not None else settings.validate_on_load
    encoding = encoding or settings.csv_encoding

    if not file_path.exists():
        msg = f"OMOP CSV file not found: {file_path}"
        logger.error(msg)
        raise FileNotFoundError(msg)

    start_time = time.time()
    logger.info(
        "Loading OMOP table: %s",
        table_name,
        extra={"file_path": str(file_path), "table_name": table_name},
    )

    try:
        df = pl.read_csv(
            file_path,
            encoding=encoding,
            null_values=["", "NULL"],
            try_parse_dates=True,
            infer_schema_length=settings.infer_schema_length,
        )

        if table_name in CDM_53_RENAMES:
            renames = {k: v for k, v in CDM_53_RENAMES[table_name].items() if k in df.columns}
            if renames:
                df = df.rename(renames)

        duration = time.time() - start_time
        logger.info(
            "Loaded %s rows in %.2fs",
            f"{len(df):,}",
            duration,
            extra={
                "table_name": table_name,
                "row_count": len(df),
                "duration_seconds": duration,
            },
        )

        if validate:
            errors = validate_omop_table_schema(df, table_name)
            if errors:
                error_summary = errors[:5]  # Log first 5 errors
                logger.warning(
                    "Validation found %d issues in %s",
                    len(errors),
                    table_name,
                    extra={
                        "table_name": table_name,
                        "error_count": len(errors),
                        "sample_errors": error_summary,
                    },
                )
                if not settings.skip_invalid_rows:
                    _raise_validation_error(table_name, len(errors))
            # Either no errors or skip_invalid_rows is True - continue to return
        # Return the DataFrame (validated or not)
        return df  # noqa: TRY300

    except Exception:
        logger.exception(
            "Failed to load OMOP table: %s",
            table_name,
            extra={"table_name": table_name, "file_path": str(file_path)},
        )
        raise

load_omop_dataset(data_dir=None, table_names=None, *, validate=True)

Load multiple OMOP tables from directory.

Parameters:

Name Type Description Default
data_dir Path | None

Directory containing OMOP CSV files. Uses config default if None.

None
table_names list[str] | None

List of table names to load. Loads all CSV files if None.

None
validate bool

Whether to validate schemas during loading.

True

Returns:

Type Description
dict[str, DataFrame]

Dictionary mapping table names to DataFrames.

Source code in src/pluginlake/omop/loader.py
def load_omop_dataset(
    data_dir: Path | None = None,
    table_names: list[str] | None = None,
    *,
    validate: bool = True,
) -> dict[str, pl.DataFrame]:
    """Load multiple OMOP tables from directory.

    Args:
        data_dir: Directory containing OMOP CSV files. Uses config default if None.
        table_names: List of table names to load. Loads all CSV files if None.
        validate: Whether to validate schemas during loading.

    Returns:
        Dictionary mapping table names to DataFrames.
    """
    settings = get_omop_settings()
    data_dir = data_dir or settings.raw_data_dir

    if not data_dir.exists():
        msg = f"OMOP data directory not found: {data_dir}"
        logger.error(msg)
        raise FileNotFoundError(msg)

    csv_files = list(data_dir.glob("*.csv"))
    if not csv_files:
        msg = f"No CSV files found in {data_dir}"
        logger.warning(msg)
        return {}

    logger.info(
        "Loading OMOP dataset from %s",
        data_dir,
        extra={"data_dir": str(data_dir), "csv_count": len(csv_files)},
    )

    tables = {}
    for csv_file in csv_files:
        table_name = csv_file.stem.lower()

        if table_names and table_name not in table_names:
            continue

        try:
            tables[table_name] = load_omop_table(csv_file, table_name, validate=validate)
        except (FileNotFoundError, ValueError, OSError):
            logger.exception(
                "Skipping table %s due to error",
                table_name,
                extra={"table_name": table_name},
            )

    logger.info(
        "Loaded %d OMOP tables",
        len(tables),
        extra={"loaded_tables": list(tables.keys())},
    )

    return tables

load_vocabulary_table(file_path, table_name, *, validate=None, encoding=None)

Load OMOP vocabulary file into validated Polars DataFrame.

Parameters:

Name Type Description Default
file_path Path

Path to vocabulary file (CSV or TSV).

required
table_name str

Vocabulary table name (e.g., 'concept', 'vocabulary').

required
validate bool | None

Run schema validation. Uses config default if None.

None
encoding str | None

File encoding. Uses config default if None.

None

Returns:

Type Description
DataFrame

Polars DataFrame with validated schema.

Raises:

Type Description
FileNotFoundError

If file doesn't exist.

ValueError

If validation fails and skip_invalid_rows is False.

Source code in src/pluginlake/omop/loader.py
def load_vocabulary_table(
    file_path: Path,
    table_name: str,
    *,
    validate: bool | None = None,
    encoding: str | None = None,
) -> pl.DataFrame:
    """Load OMOP vocabulary file into validated Polars DataFrame.

    Args:
        file_path: Path to vocabulary file (CSV or TSV).
        table_name: Vocabulary table name (e.g., 'concept', 'vocabulary').
        validate: Run schema validation. Uses config default if None.
        encoding: File encoding. Uses config default if None.

    Returns:
        Polars DataFrame with validated schema.

    Raises:
        FileNotFoundError: If file doesn't exist.
        ValueError: If validation fails and skip_invalid_rows is False.
    """
    settings = get_omop_settings()
    validate = validate if validate is not None else settings.validate_on_load
    encoding = encoding or settings.csv_encoding

    if not file_path.exists():
        msg = f"Vocabulary file not found: {file_path}"
        logger.error(msg, extra={"file_path": str(file_path), "table_name": table_name})
        raise FileNotFoundError(msg)

    start_time = time.time()
    logger.info(
        "Loading vocabulary table: %s",
        table_name,
        extra={"file_path": str(file_path), "table_name": table_name},
    )

    try:
        separator = _detect_separator(file_path, encoding)
        df = pl.read_csv(
            file_path,
            encoding=encoding,
            separator=separator,
            quote_char=None,
            null_values=["", "NULL"],
            try_parse_dates=True,
            infer_schema_length=settings.infer_schema_length,
        )

        duration = time.time() - start_time
        file_size_mb = file_path.stat().st_size / (1024 * 1024)
        logger.info(
            "Loaded %s rows (%.2f MB) in %.2fs",
            f"{len(df):,}",
            file_size_mb,
            duration,
            extra={
                "table_name": table_name,
                "row_count": len(df),
                "file_size_mb": file_size_mb,
                "duration_seconds": duration,
            },
        )

        if validate:
            errors = validate_vocabulary_table_schema(df, table_name)
            if errors:
                error_summary = errors[:5]
                logger.warning(
                    "Validation found %d issues in vocabulary %s",
                    len(errors),
                    table_name,
                    extra={
                        "table_name": table_name,
                        "error_count": len(errors),
                        "sample_errors": error_summary,
                    },
                )
                if not settings.skip_invalid_rows:
                    _raise_validation_error(table_name, len(errors))

        return df  # noqa: TRY300

    except Exception:
        logger.exception(
            "Failed to load vocabulary table: %s",
            table_name,
            extra={"table_name": table_name, "file_path": str(file_path)},
        )
        raise

load_vocabulary_dataset(data_dir=None, table_names=None, *, validate=True)

Load OMOP vocabulary tables from directory.

Parameters:

Name Type Description Default
data_dir Path | None

Directory containing vocabulary files. Uses config default if None.

None
table_names list[str] | None

List of vocabulary table names to load. Loads all if None.

None
validate bool

Whether to validate schemas during loading.

True

Returns:

Type Description
dict[str, DataFrame]

Dictionary mapping table names to DataFrames.

Source code in src/pluginlake/omop/loader.py
def load_vocabulary_dataset(
    data_dir: Path | None = None,
    table_names: list[str] | None = None,
    *,
    validate: bool = True,
) -> dict[str, pl.DataFrame]:
    """Load OMOP vocabulary tables from directory.

    Args:
        data_dir: Directory containing vocabulary files. Uses config default if None.
        table_names: List of vocabulary table names to load. Loads all if None.
        validate: Whether to validate schemas during loading.

    Returns:
        Dictionary mapping table names to DataFrames.
    """
    settings = get_omop_settings()
    data_dir = data_dir or settings.vocabulary_dir

    if not data_dir.exists():
        msg = f"Vocabulary directory not found: {data_dir}"
        logger.warning(msg, extra={"data_dir": str(data_dir)})
        return {}

    vocabulary_files = VOCABULARY_FILE_MAPPING
    if table_names:
        include = set(table_names)
        if "concept" in include:
            include.add("concept_cpt4")
        vocabulary_files = {k: v for k, v in vocabulary_files.items() if k in include}

    logger.info(
        "Loading OMOP vocabularies from %s",
        data_dir,
        extra={"data_dir": str(data_dir), "table_count": len(vocabulary_files)},
    )

    tables = {}
    for table_name, file_variants in vocabulary_files.items():
        file_path = _find_vocabulary_file(data_dir, table_name, file_variants)
        if not file_path:
            continue

        try:
            tables[table_name] = load_vocabulary_table(file_path, table_name, validate=validate)
        except (FileNotFoundError, ValueError, OSError):
            logger.exception(
                "Skipping vocabulary table %s due to error",
                table_name,
                extra={"table_name": table_name},
            )

    if "concept_cpt4" in tables:
        if "concept" in tables:
            tables["concept"] = pl.concat([tables["concept"], tables.pop("concept_cpt4")])
            logger.info("Merged concept_cpt4 into concept.")
        else:
            tables.pop("concept_cpt4")

    logger.info(
        "Loaded %d vocabulary tables",
        len(tables),
        extra={"loaded_tables": list(tables.keys())},
    )

    return tables

Queries

pluginlake.omop.queries

High-level query API for OMOP data.

Provides ergonomic Python functions for common OMOP analytical queries.

get_persons(con=None, *, gender_concept_id=None, year_of_birth_min=None, year_of_birth_max=None, race_concept_id=None, ethnicity_concept_id=None, limit=None)

Query person records with optional filters.

Source code in src/pluginlake/omop/queries.py
def get_persons(
    con: duckdb.DuckDBPyConnection | None = None,
    *,
    gender_concept_id: int | None = None,
    year_of_birth_min: int | None = None,
    year_of_birth_max: int | None = None,
    race_concept_id: int | None = None,
    ethnicity_concept_id: int | None = None,
    limit: int | None = None,
) -> pl.DataFrame:
    """Query person records with optional filters."""
    conditions = []
    params = []

    add_filter(conditions, params, "gender_concept_id", "=", gender_concept_id)
    add_filter(conditions, params, "year_of_birth", ">=", year_of_birth_min)
    add_filter(conditions, params, "year_of_birth", "<=", year_of_birth_max)
    add_filter(conditions, params, "race_concept_id", "=", race_concept_id)
    add_filter(conditions, params, "ethnicity_concept_id", "=", ethnicity_concept_id)

    query, params = build_filter_query("ducklake.omop.person", conditions, params, limit=limit)

    with ensure_connection(con) as conn:
        return execute_query(
            conn,
            query,
            params,
            lambda: logger.info("Querying persons", extra={"filter_count": len(conditions)}),
        )

get_conditions_for_person(person_id, con=None, *, condition_concept_id=None, start_date=None, end_date=None)

Query condition occurrences for a person with optional filters.

Source code in src/pluginlake/omop/queries.py
def get_conditions_for_person(
    person_id: int,
    con: duckdb.DuckDBPyConnection | None = None,
    *,
    condition_concept_id: int | None = None,
    start_date: date | None = None,
    end_date: date | None = None,
) -> pl.DataFrame:
    """Query condition occurrences for a person with optional filters."""
    conditions = []
    params = []

    add_filter(conditions, params, "person_id", "=", person_id)
    add_filter(conditions, params, "condition_concept_id", "=", condition_concept_id)
    add_filter(conditions, params, "condition_start_date", ">=", start_date)
    add_filter(conditions, params, "condition_start_date", "<=", end_date)

    query, params = build_filter_query("ducklake.omop.condition_occurrence", conditions, params)

    with ensure_connection(con) as conn:
        return execute_query(
            conn,
            query,
            params,
            lambda: logger.info(
                "Querying conditions for person",
                extra={"person_id": person_id, "filter_count": len(conditions)},
            ),
        )

get_observations_for_person(person_id, con=None, *, observation_concept_id=None, start_date=None, end_date=None)

Query observations for a person with optional filters.

Source code in src/pluginlake/omop/queries.py
def get_observations_for_person(
    person_id: int,
    con: duckdb.DuckDBPyConnection | None = None,
    *,
    observation_concept_id: int | None = None,
    start_date: date | None = None,
    end_date: date | None = None,
) -> pl.DataFrame:
    """Query observations for a person with optional filters."""
    conditions = []
    params = []

    add_filter(conditions, params, "person_id", "=", person_id)
    add_filter(conditions, params, "observation_concept_id", "=", observation_concept_id)
    add_filter(conditions, params, "observation_date", ">=", start_date)
    add_filter(conditions, params, "observation_date", "<=", end_date)

    query, params = build_filter_query("ducklake.omop.observation", conditions, params)

    with ensure_connection(con) as conn:
        return execute_query(
            conn,
            query,
            params,
            lambda: logger.info(
                "Querying observations for person",
                extra={"person_id": person_id, "filter_count": len(conditions)},
            ),
        )

get_visits_for_person(person_id, con=None, *, visit_concept_id=None, start_date=None, end_date=None)

Query visit occurrences for a person with optional filters.

Source code in src/pluginlake/omop/queries.py
def get_visits_for_person(
    person_id: int,
    con: duckdb.DuckDBPyConnection | None = None,
    *,
    visit_concept_id: int | None = None,
    start_date: date | None = None,
    end_date: date | None = None,
) -> pl.DataFrame:
    """Query visit occurrences for a person with optional filters."""
    conditions = []
    params = []

    add_filter(conditions, params, "person_id", "=", person_id)
    add_filter(conditions, params, "visit_concept_id", "=", visit_concept_id)
    add_filter(conditions, params, "visit_start_date", ">=", start_date)
    add_filter(conditions, params, "visit_start_date", "<=", end_date)

    query, params = build_filter_query("ducklake.omop.visit_occurrence", conditions, params)

    with ensure_connection(con) as conn:
        return execute_query(
            conn,
            query,
            params,
            lambda: logger.info(
                "Querying visits for person",
                extra={"person_id": person_id, "filter_count": len(conditions)},
            ),
        )

get_cohort(con=None, *, has_condition_concept_id=None, has_drug_concept_id=None, min_age=None, max_age=None, gender_concept_id=None)

Select a cohort of persons matching the given clinical criteria.

Source code in src/pluginlake/omop/queries.py
def get_cohort(
    con: duckdb.DuckDBPyConnection | None = None,
    *,
    has_condition_concept_id: int | None = None,
    has_drug_concept_id: int | None = None,
    min_age: int | None = None,
    max_age: int | None = None,
    gender_concept_id: int | None = None,
) -> pl.DataFrame:
    """Select a cohort of persons matching the given clinical criteria."""
    from_clause = "ducklake.omop.person p"
    where_clauses = []
    params = []

    if has_condition_concept_id is not None:
        from_clause += """
            INNER JOIN ducklake.omop.condition_occurrence co
                ON p.person_id = co.person_id
        """
        add_filter(where_clauses, params, "co.condition_concept_id", "=", has_condition_concept_id)

    if has_drug_concept_id is not None:
        from_clause += """
            INNER JOIN ducklake.omop.drug_exposure de
                ON p.person_id = de.person_id
        """
        add_filter(where_clauses, params, "de.drug_concept_id", "=", has_drug_concept_id)

    current_year = datetime.now(UTC).date().year
    if min_age is not None:
        max_birth_year = current_year - min_age
        add_filter(where_clauses, params, "p.year_of_birth", "<=", max_birth_year)

    if max_age is not None:
        min_birth_year = current_year - max_age
        add_filter(where_clauses, params, "p.year_of_birth", ">=", min_birth_year)

    add_filter(where_clauses, params, "p.gender_concept_id", "=", gender_concept_id)

    where_clause = " AND ".join(where_clauses) if where_clauses else "1=1"
    query = "SELECT DISTINCT p.* FROM " + from_clause + " WHERE " + where_clause

    with ensure_connection(con) as conn:
        return execute_query(
            conn,
            query,
            params,
            lambda: logger.info(
                "Selecting cohort",
                extra={
                    "criteria_count": len(where_clauses),
                    "has_condition": has_condition_concept_id is not None,
                    "has_drug": has_drug_concept_id is not None,
                },
            ),
        )

get_measurement_values(person_id, measurement_concept_id, con=None, *, start_date=None, end_date=None)

Query measurement values for a person and concept with optional date filters.

Source code in src/pluginlake/omop/queries.py
def get_measurement_values(
    person_id: int,
    measurement_concept_id: int,
    con: duckdb.DuckDBPyConnection | None = None,
    *,
    start_date: date | None = None,
    end_date: date | None = None,
) -> pl.DataFrame:
    """Query measurement values for a person and concept with optional date filters."""
    conditions = []
    params = []

    add_filter(conditions, params, "person_id", "=", person_id)
    add_filter(conditions, params, "measurement_concept_id", "=", measurement_concept_id)
    add_filter(conditions, params, "measurement_date", ">=", start_date)
    add_filter(conditions, params, "measurement_date", "<=", end_date)

    query, params = build_filter_query("ducklake.omop.measurement", conditions, params, order_by="measurement_date")

    with ensure_connection(con) as conn:
        return execute_query(
            conn,
            query,
            params,
            lambda: logger.info(
                "Querying measurement values",
                extra={
                    "person_id": person_id,
                    "measurement_concept_id": measurement_concept_id,
                },
            ),
        )

Schemas

pluginlake.omop.schemas

OMOP CDM table schemas.

Pydantic models for OMOP CDM v5.4 tables. Used for validation and type checking.

Person

Bases: BaseModel

OMOP Person table schema.

Source code in src/pluginlake/omop/schemas.py
class Person(BaseModel):
    """OMOP Person table schema."""

    person_id: int = Field(description="Unique person identifier")
    gender_concept_id: int = Field(description="Gender concept from vocabulary")
    year_of_birth: int = Field(description="Year of birth")
    month_of_birth: int | None = Field(default=None, description="Month of birth")
    day_of_birth: int | None = Field(default=None, description="Day of birth")
    birth_datetime: datetime | None = Field(default=None, description="Birth datetime")
    race_concept_id: int = Field(description="Race concept from vocabulary")
    ethnicity_concept_id: int = Field(description="Ethnicity concept from vocabulary")
    location_id: int | None = Field(default=None, description="Location foreign key")
    provider_id: int | None = Field(default=None, description="Provider foreign key")
    care_site_id: int | None = Field(default=None, description="Care site foreign key")
    person_source_value: str | None = Field(default=None, description="Source person ID")
    gender_source_value: str | None = Field(default=None, description="Source gender value")
    gender_source_concept_id: int | None = Field(default=None, description="Source gender concept")
    race_source_value: str | None = Field(default=None, description="Source race value")
    race_source_concept_id: int | None = Field(default=None, description="Source race concept")
    ethnicity_source_value: str | None = Field(default=None, description="Source ethnicity value")
    ethnicity_source_concept_id: int | None = Field(default=None, description="Source ethnicity concept")

ObservationPeriod

Bases: BaseModel

OMOP Observation Period table schema.

Source code in src/pluginlake/omop/schemas.py
class ObservationPeriod(BaseModel):
    """OMOP Observation Period table schema."""

    observation_period_id: int = Field(description="Unique observation period identifier")
    person_id: int = Field(description="Person foreign key")
    observation_period_start_date: date = Field(description="Start date of observation period")
    observation_period_end_date: date = Field(description="End date of observation period")
    period_type_concept_id: int = Field(description="Period type concept")

VisitOccurrence

Bases: BaseModel

OMOP Visit Occurrence table schema.

Source code in src/pluginlake/omop/schemas.py
class VisitOccurrence(BaseModel):
    """OMOP Visit Occurrence table schema."""

    visit_occurrence_id: int = Field(description="Unique visit occurrence identifier")
    person_id: int = Field(description="Person foreign key")
    visit_concept_id: int = Field(description="Visit concept from vocabulary")
    visit_start_date: date = Field(description="Visit start date")
    visit_start_datetime: datetime | None = Field(default=None, description="Visit start datetime")
    visit_end_date: date = Field(description="Visit end date")
    visit_end_datetime: datetime | None = Field(default=None, description="Visit end datetime")
    visit_type_concept_id: int = Field(description="Visit type concept")
    provider_id: int | None = Field(default=None, description="Provider foreign key")
    care_site_id: int | None = Field(default=None, description="Care site foreign key")
    visit_source_value: str | None = Field(default=None, description="Source visit value")
    visit_source_concept_id: int | None = Field(default=None, description="Source visit concept")
    admitted_from_concept_id: int | None = Field(default=None, description="Admitted from concept")
    admitted_from_source_value: str | None = Field(default=None, description="Admitted from source value")
    discharged_to_concept_id: int | None = Field(default=None, description="Discharged to concept")
    discharged_to_source_value: str | None = Field(default=None, description="Discharged to source value")
    preceding_visit_occurrence_id: int | None = Field(default=None, description="Preceding visit foreign key")

VisitDetail

Bases: BaseModel

OMOP Visit Detail table schema.

Source code in src/pluginlake/omop/schemas.py
class VisitDetail(BaseModel):
    """OMOP Visit Detail table schema."""

    visit_detail_id: int = Field(description="Unique visit detail identifier")
    person_id: int = Field(description="Person foreign key")
    visit_detail_concept_id: int = Field(description="Visit detail concept from vocabulary")
    visit_detail_start_date: date = Field(description="Visit detail start date")
    visit_detail_start_datetime: datetime | None = Field(default=None, description="Visit detail start datetime")
    visit_detail_end_date: date = Field(description="Visit detail end date")
    visit_detail_end_datetime: datetime | None = Field(default=None, description="Visit detail end datetime")
    visit_detail_type_concept_id: int = Field(description="Visit detail type concept")
    provider_id: int | None = Field(default=None, description="Provider foreign key")
    care_site_id: int | None = Field(default=None, description="Care site foreign key")
    visit_detail_source_value: str | None = Field(default=None, description="Source visit detail value")
    visit_detail_source_concept_id: int | None = Field(default=None, description="Source visit detail concept")
    admitted_from_concept_id: int | None = Field(default=None, description="Admitted from concept")
    admitted_from_source_value: str | None = Field(default=None, description="Admitted from source value")
    discharged_to_source_value: str | None = Field(default=None, description="Discharged to source value")
    discharged_to_concept_id: int | None = Field(default=None, description="Discharged to concept")
    preceding_visit_detail_id: int | None = Field(default=None, description="Preceding visit detail foreign key")
    parent_visit_detail_id: int | None = Field(default=None, description="Parent visit detail foreign key")
    visit_occurrence_id: int = Field(description="Visit occurrence foreign key")

ConditionOccurrence

Bases: BaseModel

OMOP Condition Occurrence table schema.

Source code in src/pluginlake/omop/schemas.py
class ConditionOccurrence(BaseModel):
    """OMOP Condition Occurrence table schema."""

    condition_occurrence_id: int = Field(description="Unique condition occurrence identifier")
    person_id: int = Field(description="Person foreign key")
    condition_concept_id: int = Field(description="Condition concept from vocabulary")
    condition_start_date: date = Field(description="Condition start date")
    condition_start_datetime: datetime | None = Field(default=None, description="Condition start datetime")
    condition_end_date: date | None = Field(default=None, description="Condition end date")
    condition_end_datetime: datetime | None = Field(default=None, description="Condition end datetime")
    condition_type_concept_id: int = Field(description="Condition type concept")
    condition_status_concept_id: int | None = Field(default=None, description="Condition status concept")
    stop_reason: str | None = Field(default=None, description="Stop reason")
    provider_id: int | None = Field(default=None, description="Provider foreign key")
    visit_occurrence_id: int | None = Field(default=None, description="Visit occurrence foreign key")
    visit_detail_id: int | None = Field(default=None, description="Visit detail foreign key")
    condition_source_value: str | None = Field(default=None, description="Source condition value")
    condition_source_concept_id: int | None = Field(default=None, description="Source condition concept")
    condition_status_source_value: str | None = Field(default=None, description="Source status value")

DrugExposure

Bases: BaseModel

OMOP Drug Exposure table schema.

Source code in src/pluginlake/omop/schemas.py
class DrugExposure(BaseModel):
    """OMOP Drug Exposure table schema."""

    drug_exposure_id: int = Field(description="Unique drug exposure identifier")
    person_id: int = Field(description="Person foreign key")
    drug_concept_id: int = Field(description="Drug concept from vocabulary")
    drug_exposure_start_date: date = Field(description="Drug exposure start date")
    drug_exposure_start_datetime: datetime | None = Field(default=None, description="Drug exposure start datetime")
    drug_exposure_end_date: date = Field(description="Drug exposure end date")
    drug_exposure_end_datetime: datetime | None = Field(default=None, description="Drug exposure end datetime")
    verbatim_end_date: date | None = Field(default=None, description="Verbatim end date")
    drug_type_concept_id: int = Field(description="Drug type concept")
    stop_reason: str | None = Field(default=None, description="Stop reason")
    refills: int | None = Field(default=None, description="Number of refills")
    quantity: float | None = Field(default=None, description="Quantity")
    days_supply: int | None = Field(default=None, description="Days supply")
    sig: str | None = Field(default=None, description="Sig")
    route_concept_id: int | None = Field(default=None, description="Route concept")
    lot_number: str | None = Field(default=None, description="Lot number")
    provider_id: int | None = Field(default=None, description="Provider foreign key")
    visit_occurrence_id: int | None = Field(default=None, description="Visit occurrence foreign key")
    visit_detail_id: int | None = Field(default=None, description="Visit detail foreign key")
    drug_source_value: str | None = Field(default=None, description="Source drug value")
    drug_source_concept_id: int | None = Field(default=None, description="Source drug concept")
    route_source_value: str | None = Field(default=None, description="Route source value")
    dose_unit_source_value: str | None = Field(default=None, description="Dose unit source value")

ProcedureOccurrence

Bases: BaseModel

OMOP Procedure Occurrence table schema.

Source code in src/pluginlake/omop/schemas.py
class ProcedureOccurrence(BaseModel):
    """OMOP Procedure Occurrence table schema."""

    procedure_occurrence_id: int = Field(description="Unique procedure occurrence identifier")
    person_id: int = Field(description="Person foreign key")
    procedure_concept_id: int = Field(description="Procedure concept from vocabulary")
    procedure_date: date = Field(description="Procedure date")
    procedure_datetime: datetime | None = Field(default=None, description="Procedure datetime")
    procedure_end_date: date | None = Field(default=None, description="Procedure end date")
    procedure_end_datetime: datetime | None = Field(default=None, description="Procedure end datetime")
    procedure_type_concept_id: int = Field(description="Procedure type concept")
    modifier_concept_id: int | None = Field(default=None, description="Modifier concept")
    quantity: int | None = Field(default=None, description="Quantity")
    provider_id: int | None = Field(default=None, description="Provider foreign key")
    visit_occurrence_id: int | None = Field(default=None, description="Visit occurrence foreign key")
    visit_detail_id: int | None = Field(default=None, description="Visit detail foreign key")
    procedure_source_value: str | None = Field(default=None, description="Source procedure value")
    procedure_source_concept_id: int | None = Field(default=None, description="Source procedure concept")
    modifier_source_value: str | None = Field(default=None, description="Modifier source value")

DeviceExposure

Bases: BaseModel

OMOP Device Exposure table schema.

Source code in src/pluginlake/omop/schemas.py
class DeviceExposure(BaseModel):
    """OMOP Device Exposure table schema."""

    device_exposure_id: int = Field(description="Unique device exposure identifier")
    person_id: int = Field(description="Person foreign key")
    device_concept_id: int = Field(description="Device concept from vocabulary")
    device_exposure_start_date: date = Field(description="Device exposure start date")
    device_exposure_start_datetime: datetime | None = Field(default=None, description="Device exposure start datetime")
    device_exposure_end_date: date | None = Field(default=None, description="Device exposure end date")
    device_exposure_end_datetime: datetime | None = Field(default=None, description="Device exposure end datetime")
    device_type_concept_id: int = Field(description="Device type concept")
    unique_device_id: str | None = Field(default=None, description="Unique device ID")
    production_id: str | None = Field(default=None, description="Production ID")
    quantity: int | None = Field(default=None, description="Quantity")
    provider_id: int | None = Field(default=None, description="Provider foreign key")
    visit_occurrence_id: int | None = Field(default=None, description="Visit occurrence foreign key")
    visit_detail_id: int | None = Field(default=None, description="Visit detail foreign key")
    device_source_value: str | None = Field(default=None, description="Source device value")
    device_source_concept_id: int | None = Field(default=None, description="Source device concept")
    unit_concept_id: int | None = Field(default=None, description="Unit concept")
    unit_source_value: str | None = Field(default=None, description="Unit source value")
    unit_source_concept_id: int | None = Field(default=None, description="Unit source concept")

Measurement

Bases: BaseModel

OMOP Measurement table schema.

Source code in src/pluginlake/omop/schemas.py
class Measurement(BaseModel):
    """OMOP Measurement table schema."""

    measurement_id: int = Field(description="Unique measurement identifier")
    person_id: int = Field(description="Person foreign key")
    measurement_concept_id: int = Field(description="Measurement concept from vocabulary")
    measurement_date: date = Field(description="Measurement date")
    measurement_datetime: datetime | None = Field(default=None, description="Measurement datetime")
    measurement_time: str | None = Field(default=None, description="Measurement time")
    measurement_type_concept_id: int = Field(description="Measurement type concept")
    operator_concept_id: int | None = Field(default=None, description="Operator concept")
    value_as_number: float | None = Field(default=None, description="Value as number")
    value_as_concept_id: int | None = Field(default=None, description="Value as concept")
    unit_concept_id: int | None = Field(default=None, description="Unit concept")
    range_low: float | None = Field(default=None, description="Range low")
    range_high: float | None = Field(default=None, description="Range high")
    provider_id: int | None = Field(default=None, description="Provider foreign key")
    visit_occurrence_id: int | None = Field(default=None, description="Visit occurrence foreign key")
    visit_detail_id: int | None = Field(default=None, description="Visit detail foreign key")
    measurement_source_value: str | None = Field(default=None, description="Source measurement value")
    measurement_source_concept_id: int | None = Field(default=None, description="Source measurement concept")
    unit_source_value: str | None = Field(default=None, description="Unit source value")
    unit_source_concept_id: int | None = Field(default=None, description="Unit source concept")
    value_source_value: str | None = Field(default=None, description="Value source value")
    measurement_event_id: int | None = Field(default=None, description="Measurement event ID")
    meas_event_field_concept_id: int | None = Field(default=None, description="Measurement event field concept")

Observation

Bases: BaseModel

OMOP Observation table schema.

Source code in src/pluginlake/omop/schemas.py
class Observation(BaseModel):
    """OMOP Observation table schema."""

    observation_id: int = Field(description="Unique observation identifier")
    person_id: int = Field(description="Person foreign key")
    observation_concept_id: int = Field(description="Observation concept from vocabulary")
    observation_date: date = Field(description="Observation date")
    observation_datetime: datetime | None = Field(default=None, description="Observation datetime")
    observation_type_concept_id: int = Field(description="Observation type concept")
    value_as_number: float | None = Field(default=None, description="Value as number")
    value_as_string: str | None = Field(default=None, description="Value as string")
    value_as_concept_id: int | None = Field(default=None, description="Value as concept")
    qualifier_concept_id: int | None = Field(default=None, description="Qualifier concept")
    unit_concept_id: int | None = Field(default=None, description="Unit concept")
    provider_id: int | None = Field(default=None, description="Provider foreign key")
    visit_occurrence_id: int | None = Field(default=None, description="Visit occurrence foreign key")
    visit_detail_id: int | None = Field(default=None, description="Visit detail foreign key")
    observation_source_value: str | None = Field(default=None, description="Source observation value")
    observation_source_concept_id: int | None = Field(default=None, description="Source observation concept")
    unit_source_value: str | None = Field(default=None, description="Unit source value")
    qualifier_source_value: str | None = Field(default=None, description="Qualifier source value")
    value_source_value: str | None = Field(default=None, description="Value source value")
    observation_event_id: int | None = Field(default=None, description="Observation event ID")
    obs_event_field_concept_id: int | None = Field(default=None, description="Observation event field concept")

Death

Bases: BaseModel

OMOP Death table schema.

Source code in src/pluginlake/omop/schemas.py
class Death(BaseModel):
    """OMOP Death table schema."""

    person_id: int = Field(description="Person foreign key")
    death_date: date = Field(description="Death date")
    death_datetime: datetime | None = Field(default=None, description="Death datetime")
    death_type_concept_id: int = Field(description="Death type concept")
    cause_concept_id: int | None = Field(default=None, description="Cause concept")
    cause_source_value: str | None = Field(default=None, description="Cause source value")
    cause_source_concept_id: int | None = Field(default=None, description="Cause source concept")

Note

Bases: BaseModel

OMOP Note table schema.

Source code in src/pluginlake/omop/schemas.py
class Note(BaseModel):
    """OMOP Note table schema."""

    note_id: int = Field(description="Unique note identifier")
    person_id: int = Field(description="Person foreign key")
    note_date: date = Field(description="Note date")
    note_datetime: datetime | None = Field(default=None, description="Note datetime")
    note_type_concept_id: int = Field(description="Note type concept")
    note_class_concept_id: int = Field(description="Note class concept")
    note_title: str | None = Field(default=None, description="Note title")
    note_text: str = Field(description="Note text")
    encoding_concept_id: int = Field(description="Encoding concept")
    language_concept_id: int = Field(description="Language concept")
    provider_id: int | None = Field(default=None, description="Provider foreign key")
    visit_occurrence_id: int | None = Field(default=None, description="Visit occurrence foreign key")
    visit_detail_id: int | None = Field(default=None, description="Visit detail foreign key")
    note_source_value: str | None = Field(default=None, description="Note source value")
    note_event_id: int | None = Field(default=None, description="Note event ID")
    note_event_field_concept_id: int | None = Field(default=None, description="Note event field concept")

NoteNlp

Bases: BaseModel

OMOP Note NLP table schema.

Source code in src/pluginlake/omop/schemas.py
class NoteNlp(BaseModel):
    """OMOP Note NLP table schema."""

    note_nlp_id: int = Field(description="Unique note NLP identifier")
    note_id: int = Field(description="Note foreign key")
    section_concept_id: int | None = Field(default=None, description="Section concept")
    snippet: str | None = Field(default=None, description="Snippet")
    offset: str | None = Field(default=None, description="Offset")
    lexical_variant: str = Field(description="Lexical variant")
    note_nlp_concept_id: int | None = Field(default=None, description="Note NLP concept")
    note_nlp_source_concept_id: int | None = Field(default=None, description="Note NLP source concept")
    nlp_system: str | None = Field(default=None, description="NLP system")
    nlp_date: date = Field(description="NLP date")
    nlp_datetime: datetime | None = Field(default=None, description="NLP datetime")
    term_exists: str | None = Field(default=None, description="Term exists")
    term_temporal: str | None = Field(default=None, description="Term temporal")
    term_modifiers: str | None = Field(default=None, description="Term modifiers")

Specimen

Bases: BaseModel

OMOP Specimen table schema.

Source code in src/pluginlake/omop/schemas.py
class Specimen(BaseModel):
    """OMOP Specimen table schema."""

    specimen_id: int = Field(description="Unique specimen identifier")
    person_id: int = Field(description="Person foreign key")
    specimen_concept_id: int = Field(description="Specimen concept from vocabulary")
    specimen_type_concept_id: int = Field(description="Specimen type concept")
    specimen_date: date = Field(description="Specimen date")
    specimen_datetime: datetime | None = Field(default=None, description="Specimen datetime")
    quantity: float | None = Field(default=None, description="Quantity")
    unit_concept_id: int | None = Field(default=None, description="Unit concept")
    anatomic_site_concept_id: int | None = Field(default=None, description="Anatomic site concept")
    disease_status_concept_id: int | None = Field(default=None, description="Disease status concept")
    specimen_source_id: str | None = Field(default=None, description="Specimen source ID")
    specimen_source_value: str | None = Field(default=None, description="Specimen source value")
    unit_source_value: str | None = Field(default=None, description="Unit source value")
    anatomic_site_source_value: str | None = Field(default=None, description="Anatomic site source value")
    disease_status_source_value: str | None = Field(default=None, description="Disease status source value")

ConditionEra

Bases: BaseModel

OMOP Condition Era table schema.

Source code in src/pluginlake/omop/schemas.py
class ConditionEra(BaseModel):
    """OMOP Condition Era table schema."""

    condition_era_id: int = Field(description="Unique condition era identifier")
    person_id: int = Field(description="Person foreign key")
    condition_concept_id: int = Field(description="Condition concept from vocabulary")
    condition_era_start_date: date = Field(description="Start date of condition era")
    condition_era_end_date: date = Field(description="End date of condition era")
    condition_occurrence_count: int | None = Field(default=None, description="Number of condition occurrences")

DrugEra

Bases: BaseModel

OMOP Drug Era table schema.

Source code in src/pluginlake/omop/schemas.py
class DrugEra(BaseModel):
    """OMOP Drug Era table schema."""

    drug_era_id: int = Field(description="Unique drug era identifier")
    person_id: int = Field(description="Person foreign key")
    drug_concept_id: int = Field(description="Drug concept from vocabulary")
    drug_era_start_date: date = Field(description="Start date of drug era")
    drug_era_end_date: date = Field(description="End date of drug era")
    drug_exposure_count: int | None = Field(default=None, description="Number of drug exposures")
    gap_days: int | None = Field(default=None, description="Gap days")

FactRelationship

Bases: BaseModel

OMOP Fact Relationship table schema.

Source code in src/pluginlake/omop/schemas.py
class FactRelationship(BaseModel):
    """OMOP Fact Relationship table schema."""

    domain_concept_id_1: int = Field(description="Domain concept 1")
    fact_id_1: int = Field(description="Fact ID 1")
    domain_concept_id_2: int = Field(description="Domain concept 2")
    fact_id_2: int = Field(description="Fact ID 2")
    relationship_concept_id: int = Field(description="Relationship concept")

get_omop_schema(table_name)

Get Pydantic schema for OMOP table.

Parameters:

Name Type Description Default
table_name str

OMOP table name (lowercase with underscores).

required

Returns:

Type Description
type[BaseModel] | None

Pydantic model class or None if not found.

Source code in src/pluginlake/omop/schemas.py
def get_omop_schema(table_name: str) -> type[BaseModel] | None:
    """Get Pydantic schema for OMOP table.

    Args:
        table_name: OMOP table name (lowercase with underscores).

    Returns:
        Pydantic model class or None if not found.
    """
    return OMOP_SCHEMAS.get(table_name)

Validation

pluginlake.omop.validation

OMOP data validation functions.

ValidationError

Represents a validation error.

Source code in src/pluginlake/omop/validation.py
class ValidationError:
    """Represents a validation error."""

    def __init__(self, column: str, message: str, row_index: int | None = None) -> None:
        """Initialize validation error.

        Args:
            column: Column name where error occurred.
            message: Error description.
            row_index: Optional row index where error occurred.
        """
        self.column = column
        self.message = message
        self.row_index = row_index

    def __repr__(self) -> str:
        """String representation of validation error."""
        if self.row_index is not None:
            return f"ValidationError(row={self.row_index}, column={self.column}, message={self.message})"
        return f"ValidationError(column={self.column}, message={self.message})"

__init__(column, message, row_index=None)

Initialize validation error.

Parameters:

Name Type Description Default
column str

Column name where error occurred.

required
message str

Error description.

required
row_index int | None

Optional row index where error occurred.

None
Source code in src/pluginlake/omop/validation.py
def __init__(self, column: str, message: str, row_index: int | None = None) -> None:
    """Initialize validation error.

    Args:
        column: Column name where error occurred.
        message: Error description.
        row_index: Optional row index where error occurred.
    """
    self.column = column
    self.message = message
    self.row_index = row_index

__repr__()

String representation of validation error.

Source code in src/pluginlake/omop/validation.py
def __repr__(self) -> str:
    """String representation of validation error."""
    if self.row_index is not None:
        return f"ValidationError(row={self.row_index}, column={self.column}, message={self.message})"
    return f"ValidationError(column={self.column}, message={self.message})"

validate_omop_table_schema(df, table_name)

Check DataFrame matches OMOP CDM schema.

Parameters:

Name Type Description Default
df DataFrame

DataFrame to validate.

required
table_name str

OMOP table name.

required

Returns:

Type Description
list[ValidationError]

List of validation errors (empty if valid).

Source code in src/pluginlake/omop/validation.py
def validate_omop_table_schema(
    df: pl.DataFrame,
    table_name: str,
) -> list[ValidationError]:
    """Check DataFrame matches OMOP CDM schema.

    Args:
        df: DataFrame to validate.
        table_name: OMOP table name.

    Returns:
        List of validation errors (empty if valid).
    """
    schema = get_omop_schema(table_name)
    if not schema:
        logger.warning("No schema definition for table: %s", table_name)
        return []

    errors = []
    schema_fields = schema.model_fields
    df_columns = set(df.columns)
    expected_columns = set(schema_fields.keys())

    missing = expected_columns - df_columns
    for col in missing:
        field = schema_fields[col]
        if field.is_required():
            errors.append(ValidationError(col, "Required column missing"))

    unexpected = df_columns - expected_columns
    errors.extend(ValidationError(col, "Unexpected column not in schema") for col in unexpected)

    for col in df_columns & expected_columns:
        field = schema_fields[col]
        df_type = df[col].dtype

        expected_type = _get_expected_polars_type(field.annotation)

        if expected_type and df_type != expected_type:
            errors.append(
                ValidationError(
                    col,
                    f"Type mismatch: expected {expected_type}, got {df_type}",
                )
            )

    return errors

validate_vocabulary_table_schema(df, table_name)

Check DataFrame matches OMOP vocabulary schema.

Parameters:

Name Type Description Default
df DataFrame

DataFrame to validate.

required
table_name str

Vocabulary table name.

required

Returns:

Type Description
list[ValidationError]

List of validation errors (empty if valid).

Source code in src/pluginlake/omop/validation.py
def validate_vocabulary_table_schema(
    df: pl.DataFrame,
    table_name: str,
) -> list[ValidationError]:
    """Check DataFrame matches OMOP vocabulary schema.

    Args:
        df: DataFrame to validate.
        table_name: Vocabulary table name.

    Returns:
        List of validation errors (empty if valid).
    """
    schema = get_vocabulary_schema(table_name)
    if not schema:
        logger.warning("No schema definition for vocabulary table: %s", table_name)
        return []

    errors = []
    schema_fields = schema.model_fields
    df_columns = set(df.columns)
    expected_columns = set(schema_fields.keys())

    missing = expected_columns - df_columns
    for col in missing:
        field = schema_fields[col]
        if field.is_required():
            errors.append(ValidationError(col, "Required column missing"))

    unexpected = df_columns - expected_columns
    errors.extend(ValidationError(col, "Unexpected column not in schema") for col in unexpected)

    for col in df_columns & expected_columns:
        field = schema_fields[col]
        df_type = df[col].dtype

        expected_type = _get_expected_polars_type(field.annotation)

        if expected_type and df_type != expected_type:
            errors.append(
                ValidationError(
                    col,
                    f"Type mismatch: expected {expected_type}, got {df_type}",
                )
            )

    return errors

Vocabulary Queries

pluginlake.omop.vocabulary_queries

Query API for OMOP vocabularies.

Functions for querying OMOP controlled vocabularies.

get_concept(concept_id, con=None, *, schema='ducklake.omop_vocab')

Retrieve a single concept by ID.

Source code in src/pluginlake/omop/vocabulary_queries.py
def get_concept(
    concept_id: int,
    con: duckdb.DuckDBPyConnection | None = None,
    *,
    schema: str = "ducklake.omop_vocab",
) -> pl.DataFrame:
    """Retrieve a single concept by ID."""
    with ensure_connection(con) as conn:
        query = f"""
            SELECT
                concept_id,
                concept_name,
                domain_id,
                vocabulary_id,
                concept_class_id,
                standard_concept,
                concept_code,
                valid_start_date,
                valid_end_date,
                invalid_reason
            FROM {schema}.concept
            WHERE concept_id = $concept_id
        """  # noqa: S608

        return execute_query(
            conn,
            query,
            {"concept_id": concept_id},
            lambda: logger.debug("Retrieving concept: %d", concept_id),
        )

search_concepts(term, domain_id=None, vocabulary_id=None, *, standard_only=True, limit=100, con=None, schema='ducklake.omop_vocab')

Search concepts by name with optional domain, vocabulary, and standard filters.

Source code in src/pluginlake/omop/vocabulary_queries.py
def search_concepts(  # noqa: PLR0913
    term: str,
    domain_id: str | None = None,
    vocabulary_id: str | None = None,
    *,
    standard_only: bool = True,
    limit: int = 100,
    con: duckdb.DuckDBPyConnection | None = None,
    schema: str = "ducklake.omop_vocab",
) -> pl.DataFrame:
    """Search concepts by name with optional domain, vocabulary, and standard filters."""
    with ensure_connection(con) as conn:
        conditions = []
        params = {"term": f"%{term}%", "limit": limit}

        conditions.append("LOWER(concept_name) LIKE LOWER($term)")

        if domain_id:
            conditions.append("domain_id = $domain_id")
            params["domain_id"] = domain_id

        if vocabulary_id:
            conditions.append("vocabulary_id = $vocabulary_id")
            params["vocabulary_id"] = vocabulary_id

        if standard_only:
            conditions.append("standard_concept = 'S'")

        conditions.append("invalid_reason IS NULL")

        where_clause = " AND ".join(conditions)

        query = f"""
            SELECT
                concept_id,
                concept_name,
                domain_id,
                vocabulary_id,
                concept_class_id,
                standard_concept,
                concept_code
            FROM {schema}.concept
            WHERE {where_clause}
            ORDER BY concept_name
            LIMIT $limit
        """  # noqa: S608

        return execute_query(
            conn,
            query,
            params,
            lambda: logger.debug("Searching concepts: '%s'", term),
        )

get_concept_descendants(ancestor_concept_id, max_levels=None, *, con=None, schema='ducklake.omop_vocab')

Return all descendant concepts of an ancestor concept.

Source code in src/pluginlake/omop/vocabulary_queries.py
def get_concept_descendants(
    ancestor_concept_id: int,
    max_levels: int | None = None,
    *,
    con: duckdb.DuckDBPyConnection | None = None,
    schema: str = "ducklake.omop_vocab",
) -> pl.DataFrame:
    """Return all descendant concepts of an ancestor concept."""
    with ensure_connection(con) as conn:
        params = {"ancestor_id": ancestor_concept_id}
        max_levels_clause = ""

        if max_levels is not None:
            max_levels_clause = "AND ca.max_levels_of_separation <= $max_levels"
            params["max_levels"] = max_levels

        query = f"""
            SELECT
                c.concept_id,
                c.concept_name,
                c.domain_id,
                c.vocabulary_id,
                c.concept_class_id,
                c.standard_concept,
                ca.min_levels_of_separation,
                ca.max_levels_of_separation
            FROM {schema}.concept_ancestor ca
            JOIN {schema}.concept c ON ca.descendant_concept_id = c.concept_id
            WHERE ca.ancestor_concept_id = $ancestor_id
              {max_levels_clause}
            ORDER BY ca.min_levels_of_separation, c.concept_name
        """  # noqa: S608

        return execute_query(
            conn,
            query,
            params,
            lambda: logger.debug("Getting descendants of concept: %d", ancestor_concept_id),
        )

get_concept_ancestors(descendant_concept_id, max_levels=None, *, con=None, schema='ducklake.omop_vocab')

Return all ancestor concepts of a descendant concept.

Source code in src/pluginlake/omop/vocabulary_queries.py
def get_concept_ancestors(
    descendant_concept_id: int,
    max_levels: int | None = None,
    *,
    con: duckdb.DuckDBPyConnection | None = None,
    schema: str = "ducklake.omop_vocab",
) -> pl.DataFrame:
    """Return all ancestor concepts of a descendant concept."""
    with ensure_connection(con) as conn:
        params = {"descendant_id": descendant_concept_id}
        max_levels_clause = ""

        if max_levels is not None:
            max_levels_clause = "AND ca.max_levels_of_separation <= $max_levels"
            params["max_levels"] = max_levels

        query = f"""
            SELECT
                c.concept_id,
                c.concept_name,
                c.domain_id,
                c.vocabulary_id,
                c.concept_class_id,
                c.standard_concept,
                ca.min_levels_of_separation,
                ca.max_levels_of_separation
            FROM {schema}.concept_ancestor ca
            JOIN {schema}.concept c ON ca.ancestor_concept_id = c.concept_id
            WHERE ca.descendant_concept_id = $descendant_id
              {max_levels_clause}
            ORDER BY ca.min_levels_of_separation, c.concept_name
        """  # noqa: S608

        return execute_query(
            conn,
            query,
            params,
            lambda: logger.debug("Getting ancestors of concept: %d", descendant_concept_id),
        )

map_source_code(source_code, source_vocabulary_id, *, con=None, schema='ducklake.omop_vocab')

Map a source code to its standard OMOP concept via the source-to-concept map.

Source code in src/pluginlake/omop/vocabulary_queries.py
def map_source_code(
    source_code: str,
    source_vocabulary_id: str,
    *,
    con: duckdb.DuckDBPyConnection | None = None,
    schema: str = "ducklake.omop_vocab",
) -> pl.DataFrame:
    """Map a source code to its standard OMOP concept via the source-to-concept map."""
    with ensure_connection(con) as conn:
        query = f"""
            SELECT
                stcm.source_code,
                stcm.source_concept_id,
                stcm.source_vocabulary_id,
                stcm.source_code_description,
                stcm.target_concept_id,
                c.concept_name as target_concept_name,
                c.domain_id as target_domain_id,
                c.vocabulary_id as target_vocabulary_id,
                c.concept_class_id as target_concept_class_id,
                c.standard_concept as target_standard_concept
            FROM {schema}.source_to_concept_map stcm
            JOIN {schema}.concept c ON stcm.target_concept_id = c.concept_id
            WHERE stcm.source_code = $source_code
              AND stcm.source_vocabulary_id = $source_vocabulary_id
              AND stcm.invalid_reason IS NULL
              AND c.invalid_reason IS NULL
            ORDER BY stcm.target_concept_id
        """  # noqa: S608

        return execute_query(
            conn,
            query,
            {"source_code": source_code, "source_vocabulary_id": source_vocabulary_id},
            lambda: logger.debug("Mapping source code: %s (%s)", source_code, source_vocabulary_id),
        )

get_vocabulary_info(vocabulary_id=None, *, con=None, schema='ducklake.omop_vocab')

Return vocabulary metadata, optionally filtered by vocabulary ID.

Source code in src/pluginlake/omop/vocabulary_queries.py
def get_vocabulary_info(
    vocabulary_id: str | None = None,
    *,
    con: duckdb.DuckDBPyConnection | None = None,
    schema: str = "ducklake.omop_vocab",
) -> pl.DataFrame:
    """Return vocabulary metadata, optionally filtered by vocabulary ID."""
    with ensure_connection(con) as conn:
        where_clause = ""
        params = {}

        if vocabulary_id:
            where_clause = "WHERE vocabulary_id = $vocabulary_id"
            params["vocabulary_id"] = vocabulary_id

        query = f"""
            SELECT
                vocabulary_id,
                vocabulary_name,
                vocabulary_reference,
                vocabulary_version,
                vocabulary_concept_id
            FROM {schema}.vocabulary
            {where_clause}
            ORDER BY vocabulary_id
        """  # noqa: S608

        return execute_query(
            conn,
            query,
            params or None,
            lambda: logger.debug("Getting vocabulary info"),
        )