Skip to content

DBT: Crushing Facilities 2003 2024

File location: s3://trase-storage/brazil/logistics/abiove/out/CRUSHING_FACILITIES_2003_2024.csv

DBT model name: crushing_facilities_2003_2024

Explore on Metabase: Full table; summary statistics

Explore dependencies/lineage: link

Relies on script: trase/data/brazil/logistics/abiove/out/CRUSHING_FACILITIES_2023_20XX.py


Description

This model was auto-generated based off .yml 'lineage' files in S3. The DBT model just raises an error; the actual script that created the data lives elsewhere. The script is located at trase/data/brazil/logistics/abiove/out/CRUSHING_FACILITIES_2023_20XX.py, an update of the earlier crushing_facilities_2003_2019. It is read directly by the brazil_soy_2023_2024_v27 supply chain (IndustrialCapacity/IndustrialCapacityFacilities in preparation.py) and by brazil_soy_supply_sheds (CrushingDemand).


Details

Column Type Description
TRASE_ID VARCHAR
YEAR VARCHAR
COMPANY VARCHAR
MUNICIPALITY VARCHAR
UF VARCHAR
GEOCODE VARCHAR
CAPACITY VARCHAR
CAPACITY_SOURCE VARCHAR
CNPJ VARCHAR
LAT VARCHAR
LONG VARCHAR
RESOLUTION VARCHAR
CAPACITY_ANNUAL VARCHAR

No data tests defined 🧐

No dependencies recorded.

import pandas as pd
import re

from pprint import pprint
from trase.tools.aws.aws_helpers_cached import get_pandas_df
from trase.tools.aws.metadata import write_csv_for_upload
from trase.tools.aws.aws_helpers import read_geojson
from trase.tools import sps
from helpers.mannual_inspections import *
import numpy as np

"""
    Config Session
"""

PATH_CRUSHING_FACILITIES_2003_2022 = (
    "brazil/logistics/jjhinrichsen/out/CRUSHING_FACILITIES_2003_2022.csv"
)

PATHS_CRUSHING_CAPACITY = {
    "2023": "brazil/logistics/jjhinrichsen/jjhinrichsen_2023.csv",  # mannualy extracted from JJ pdfs
    "2024": "brazil/logistics/jjhinrichsen/jjhinrichsen_2024.csv",  # mannualy extracted from JJ pdfs
}

PATHS_CRUSHING_CAPACITY_ABIOVE = {
    "2023": "brazil/logistics/abiove/ori/Pesquisa-de-Capacidade-Instalada_2024.xlsx",  # The year 2023 is inside 2024's file
    "2024": "brazil/logistics/abiove/ori/Pesquisa-de-Capacidade-Instalada_2024.xlsx",
}

PATH_OUTPUT = "brazil/logistics/abiove/out/CRUSHING_FACILITIES_2003_2024.csv"

PATH_UF = "brazil/metadata/UF.csv"
PATH_MUNICIPALITIES = (
    "brazil/spatial/boundaries/ibge/2023/br_municipalities_wgs84_2023.geojson"
)

CAPACITY_RAW_COLS = [
    "Company",
    "location",
    "Department",
    "Capacity",
    "Process",
    "Seed",
    "Remarks",
]
CAPACITY_NEW_COLS = [
    "Company",
    "municipality",
    "uf",
    "Capacity",
    "Process",
    "Seed",
    "Remarks",
]

ABIOVE_RAW_COLS = ["Empresas", "Município", "UF", "Região", "Soja"]
ABIOVE_NEW_COLS = ["COMPANY", "MUNICIPALITY", "UF", "REGION", "PROCESS_SOY"]


"""
    Helper Functions
"""


def normalize_str(d: pd.DataFrame, col: str, clean=False):
    """
    Adjust column value characters encoding to UTF-8 and uppercase them.

    Args:
        d (pandas DataFrame): Dataframe to lookup
        col (str): String column
        clean: remove specific characters

    Returns:
        pandas DataFrame
    """
    d[col] = (
        d[col]
        .str.normalize("NFKD")
        .str.encode("ascii", errors="ignore")
        .str.decode("utf-8")
        .str.strip()
        .str.replace(r"\s+", " ", regex=True)
    )

    d[col] = d[col].str.upper()

    if clean is True:
        d[col] = (
            d[col]
            .str.replace(".", "")
            .str.replace("-", "")
            .str.replace("/", "")
            .str.replace(",", "")
            .str.replace('"', "")
            .str.split(" ")
            .str.join("")
        )

    else:
        d[col] = d[col]

    return d


"""
    Main Function
"""


def main():

    # Load input data
    df_crushing_facilities_2003_2022 = get_pandas_df(PATH_CRUSHING_FACILITIES_2003_2022)

    df_uf = get_pandas_df(PATH_UF, sep=",")

    df_crushing_facilities_2003_2022 = normalize_str(
        df_crushing_facilities_2003_2022, col="COMPANY"
    )

    df_municipalities = read_geojson(key=PATH_MUNICIPALITIES, bucket="trase-storage")

    # Link capacity to facilities in each year
    facilities_crushing_list = []
    for year, path in PATHS_CRUSHING_CAPACITY.items():

        print(f"Processing year {year}...")

        # Loading Session
        # =========================================================================================
        # Load capacity from JJ dataset (manually extracted from pdf)
        df_capacity = get_pandas_df(
            key=path, bucket="trase-storage", sep=";", usecols=CAPACITY_RAW_COLS
        )

        # Load facilities from ABIOVE
        df_abiove_facilities = get_pandas_df(
            key=PATHS_CRUSHING_CAPACITY_ABIOVE.get(year),
            bucket="trase-storage",
            xlsx=True,
            sheet_name="3.Unidades de Processamento",
            usecols="B:E,F:G,I",
            header=7,
            skipfooter=7,
        )[ABIOVE_RAW_COLS + [int(year)]]

        # Cleaning Session
        # =========================================================================================

        # 0. Renaming columns
        df_abiove_facilities = df_abiove_facilities.rename(
            columns=dict(
                zip(ABIOVE_RAW_COLS + [int(year)], ABIOVE_NEW_COLS + ["STATUS"])
            )
        )

        df_capacity = df_capacity.rename(
            columns=dict(
                zip(CAPACITY_RAW_COLS + [int(year)], CAPACITY_NEW_COLS + ["STATUS"])
            )
        )

        df_capacity.columns = [x.upper() for x in df_capacity.columns]

        # 1. Remove NAN rows
        df_abiove_facilities = df_abiove_facilities.dropna(how="all")
        df_capacity = df_capacity.dropna(how="all")

        df_abiove_facilities["YEAR"] = int(year)
        df_capacity["YEAR"] = int(year)

        # 2. Cleaning columns
        cols_to_clean = ["COMPANY", "MUNICIPALITY", "UF"]
        for col in cols_to_clean:
            df_capacity = normalize_str(df_capacity, col=col)
            df_abiove_facilities = normalize_str(df_abiove_facilities, col=col)

        # 3. Remove duplicates
        df_abiove_facilities = df_abiove_facilities.drop_duplicates()
        df_capacity = df_capacity.drop_duplicates()

        # 4. Remap values
        df_capacity = df_capacity.replace(
            {
                "MUNICIPALITY": MUNICIPALITY_REMAP,
                "UF": STATE_REMAP,
                "COMPANY": COMPANIES_REMAP,
            }
        ).rename(columns={"UF": "NO_UF"})

        df_abiove_facilities = df_abiove_facilities.replace(
            {"MUNICIPALITY": MUNICIPALITY_REMAP, "COMPANY": COMPANIES_REMAP}
        )

        # Bring state name to jj file
        df_capacity = df_capacity.merge(
            df_uf[["NO_UF", "UF", "CO_UF"]], on=["NO_UF"], how="left"
        )

        df_abiove_facilities = df_abiove_facilities.merge(
            df_uf[["NO_UF", "UF", "CO_UF"]], on=["UF"], how="left"
        )

        # Filter facilities that process soy and are active
        # =========================================================================================

        df_abiove_facilities = df_abiove_facilities.query(
            'PROCESS_SOY == "X" and STATUS == "Ativa"'
        )

        # Combining JJ Capacity and Abiove
        # =========================================================================================

        df_crushing_facilities_year = df_abiove_facilities.merge(
            df_capacity,
            how="outer",
            on=["COMPANY", "MUNICIPALITY", "UF", "YEAR", "CO_UF"],
            indicator=True,
        )[["YEAR", "COMPANY", "MUNICIPALITY", "UF", "CAPACITY", "CO_UF", "_merge"]]

        # Fix one especific company name
        df_crushing_facilities_year = df_crushing_facilities_year.replace(
            {"COMPANY": {"CNH": "GRANOL"}}
        )

        # Adjust State names
        df_correct_mun_states = pd.DataFrame(CORRECT_MUNICIPALITIES_STATES)
        df_crushing_facilities_year = df_crushing_facilities_year.merge(
            df_correct_mun_states[["MUNICIPALITY", "CO_UF", "UF_CODE"]],
            how="left",
            on=["MUNICIPALITY", "CO_UF"],
        )
        df_crushing_facilities_year.loc[
            df_crushing_facilities_year["UF_CODE"].notnull(), "CO_UF"
        ] = df_crushing_facilities_year["UF_CODE"]

        # Adding Missing Capacities
        # =========================================================================================
        df_missing_capacities = pd.DataFrame(MISSING_CAPACITIES)
        df_missing_capacities_year = df_missing_capacities.query(f"YEAR == {int(year)}")

        df_crushing_facilities_year = df_crushing_facilities_year.merge(
            df_missing_capacities_year.rename(columns={"CAPACITY": "CAPACITY_MISSING"}),
            how="left",
            on=["COMPANY", "MUNICIPALITY", "UF", "YEAR"],
        )

        df_crushing_facilities_year["CAPACITY"] = df_crushing_facilities_year[
            "CAPACITY"
        ].fillna(df_crushing_facilities_year["CAPACITY_MISSING"])
        df_crushing_facilities_year["CAPACITY_SOURCE"] = df_crushing_facilities_year[
            "CAPACITY_SOURCE"
        ].fillna("JJ")

        df_crushing_facilities_year = df_crushing_facilities_year.drop(
            columns=["_merge", "CAPACITY_MISSING"]
        )

        df_crushing_facilities_year["RESOLUTION"] = "DISTRICT"

        # Bringing additional fields (GEOCODE, CNPJ)
        # =========================================================================================

        # 1. CNPJ
        # 1.1 Finding CNPJ from historical data
        historical_data = df_crushing_facilities_2003_2022[
            ["COMPANY", "CNPJ", "MUNICIPALITY", "GEOCODE"]
        ].drop_duplicates()

        df_crushing_facilities_year = df_crushing_facilities_year.merge(
            historical_data, how="left", on=["COMPANY", "MUNICIPALITY"]
        )

        # 1.2 Finding CNPJ from mannual inspected data
        df_missing_cnpjs = pd.DataFrame(CNPJS_MISSING)
        df_crushing_facilities_year = df_crushing_facilities_year.merge(
            df_missing_cnpjs, how="left", on=["COMPANY", "MUNICIPALITY"]
        )

        df_crushing_facilities_year["CNPJ"] = np.where(
            df_crushing_facilities_year["CNPJ_x"].notnull(),
            df_crushing_facilities_year["CNPJ_x"],
            df_crushing_facilities_year["CNPJ_y"],
        )

        df_crushing_facilities_year = df_crushing_facilities_year.drop(
            columns=["CNPJ_x", "CNPJ_y"]
        )

        # Bringing trase id
        # normalize municipality name
        df_municipalities = normalize_str(df_municipalities, col="name")
        df_crushing_facilities_year = normalize_str(
            df_crushing_facilities_year, col="MUNICIPALITY"
        )
        df_crushing_facilities_year = normalize_str(
            df_crushing_facilities_year, col="CNPJ", clean=True
        )

        df_municipalities["ibge_state_id"] = df_municipalities["ibge_state_id"].astype(
            int
        )
        df_crushing_facilities_year["CO_UF"] = (
            df_crushing_facilities_year["CO_UF"].fillna(-1).astype(int)
        )

        df_crushing_facilities_year = df_crushing_facilities_year.merge(
            df_municipalities,
            how="left",
            left_on=["CO_UF", "MUNICIPALITY"],
            right_on=["ibge_state_id", "name"],
        )

        df_crushing_facilities_year["GEOCODE"] = np.where(
            df_crushing_facilities_year["GEOCODE"].isna(),
            df_crushing_facilities_year["ibge_municipality_id"],
            df_crushing_facilities_year["GEOCODE"],
        )
        df_crushing_facilities_year["LAT"] = None
        df_crushing_facilities_year["LONG"] = None

        df_crushing_facilities_year.columns = [
            x.upper() for x in df_crushing_facilities_year.columns
        ]

        df_crushing_facilities_year = df_crushing_facilities_year[
            [
                "TRASE_ID",
                "YEAR",
                "COMPANY",
                "MUNICIPALITY",
                "UF",
                "GEOCODE",
                "CAPACITY",
                "CAPACITY_SOURCE",
                "CNPJ",
                "LAT",
                "LONG",
                "RESOLUTION",
            ]
        ]

        mask = (
            df_crushing_facilities_year["TRASE_ID"].isna()
            & df_crushing_facilities_year["GEOCODE"].notna()
        )

        df_crushing_facilities_year.loc[mask, "TRASE_ID"] = (
            "BR-"
            + df_crushing_facilities_year.loc[mask, "GEOCODE"].astype(int).astype(str)
        )

        # Filter colums and add to list
        facilities_crushing_list.append(df_crushing_facilities_year)

    # Stitch together and add yearly capacity
    # =========================================================================================
    df_crushing_facilities_2003_2022["TRASE_ID"] = (
        "BR-" + df_crushing_facilities_2003_2022["GEOCODE"].fillna("").astype(str)
    )

    df_crushing_facilities_2003_2022 = df_crushing_facilities_2003_2022[
        [
            "TRASE_ID",
            "YEAR",
            "COMPANY",
            "MUNICIPALITY",
            "UF",
            "GEOCODE",
            "CAPACITY",
            "CAPACITY_SOURCE",
            "CNPJ",
            "LAT",
            "LONG",
            "RESOLUTION",
        ]
    ]

    df_crushing_facilities_output = sps.concat(
        [df_crushing_facilities_2003_2022] + facilities_crushing_list
    )

    df_crushing_facilities_output = df_crushing_facilities_output.drop_duplicates()

    # Annualizing capacity using a 288-day operational factor.
    # This accounts for ~18% downtime (65 days/year) covering scheduled maintenance,
    # soybean seasonality, and logistical constraints.
    df_crushing_facilities_output["CAPACITY_ANNUAL"] = df_crushing_facilities_output[
        "CAPACITY"
    ].mul(288)

    # Clean CNPJ
    # =========================================================================================
    df_crushing_facilities_output["CNPJ"] = df_crushing_facilities_output[
        "CNPJ"
    ].str.replace(r"\D", "", regex=True)

    # It can actually contain null CNPJs
    # assert df_crushing_facilities_output.query('YEAR > 2021 and CNPJ.isna()').shape[0] == 0, 'Contains NaN CNPJ'
    assert (
        df_crushing_facilities_output.query("YEAR > 2021 and TRASE_ID.isna()").shape[0]
        == 0
    ), "Contains NaN TRASE_ID"

    # Export
    # =========================================================================================
    write_csv_for_upload(df_crushing_facilities_output, key=PATH_OUTPUT, sep=";")


"""
    Helper Functions
"""

if __name__ == "__main__":
    main()
import pandas as pd


def model(dbt, cursor):
    raise NotImplementedError()
    return pd.DataFrame({"hello": ["world"]})