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 🧐
Models
crushing_supply_shedshubs_crushing_capacitiesseipcs_brazil_soy_partial_branches_v27silo_capacitiessilo_supply_shedssolution_domestic_unknownssolution_exports
Exposures
brazil_soy_v27_comparison_to_secex_exports(view exposure, view documentation)
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"]})
-
Dbt path:
memory.main.crushing_facilities_2003_2024 -
Containing yaml link: trase/data_pipeline/models/brazil/logistics/abiove/out/_schema.yml
-
Model file: trase/data_pipeline/models/brazil/logistics/abiove/out/crushing_facilities_2003_2024.py
-
Calls script:
trase/data/brazil/logistics/abiove/out/CRUSHING_FACILITIES_2023_20XX.py -
Tags:
mock_model,abiove,brazil,logistics,out