Post processing
View or edit on GitHub
This page is synchronized from trase/models/indonesia/palm_oil/post_processing.ipynb. Last modified on 2026-08-05 15:56 CEST by Harry Biddle.
Please view or edit the original file there; changes should be reflected here after a midnight build (CET time),
or manually triggering it with a GitHub action (link).
library(dplyr)
library(readr)
library(tidylog)
library(arrow)
library(dlookr)
library(aws.s3)
Attaching package: ‘dplyr’
The following objects are masked from ‘package:stats’:
filter, lag
The following objects are masked from ‘package:base’:
intersect, setdiff, setequal, union
Error in library(tidylog): there is no package called ‘tidylog’
Traceback:
1. stop(packageNotFoundError(package, lib.loc, sys.call()))
bucket <- "trase-storage"
region <- "eu-west-1"
Sys.setenv("AWS_DEFAULT_REGION" = region)
year = 2024
version = "v1.3"
expected_value_columns <- c("fob", "product_vol", "vol")
message("Read country lookup")
df_countries <- s3read_using(
FUN = arrow::read_parquet,
object = "postgres_views/postgres_countries.parquet",
bucket = bucket) %>%
as.data.frame()
if (!"country_id" %in% names(df_countries)) {
df_countries$country_id <- df_countries$iso_numeric_code}
df_countries <- df_countries %>%
transmute(
country_join = toupper(trimws(country_name)),
country_id = as.character(country_id),
country_trase_id = as.character(country_trase_id)) %>%
distinct(country_join, .keep_all = TRUE)
message("Read EU economic bloc lookup")
df_eu <- s3read_using(
FUN = arrow::read_parquet,
object = "postgres_views/postgres_eu_economic_blocs.parquet",
bucket = bucket) %>%
as.data.frame() %>%
transmute(
country_trase_id = as.character(country_trase_id),
economic_bloc_trase_id = "EU") %>%
distinct(country_trase_id, .keep_all = TRUE)
#process_year <- function(year, version = "v1.3") {
input_key <- paste0(
"indonesia/palm_oil/sei_pcs/", version, "/",
"results_", year, "_KM_0_MGD_30.csv"
)
output_key <- paste0(
"indonesia/palm_oil/sei_pcs/", version, "/",
"SEIPCS_INDONESIA_PALM_OIL_", year, "_KM_0_MGD_30.parquet"
)
message("Read data")
df <- s3read_using(
read_delim,
object = input_key,
bucket = bucket,
delim = ";",
show_col_types = FALSE
)
print(names(df))
print(glimpse(df))
message("Prepare results")
df <- df %>%
mutate(year = as.integer(year))
for (col in intersect(expected_value_columns, names(df))) {
df[[col]] <- as.numeric(df[[col]])
}
message("Join country ids and economic blocs")
df <- df %>%
mutate(country = ifelse(exporter == "DOMESTIC PROCESSING AND CONSUMPTION", "INDONESIA", country)) %>%
left_join(
df_countries %>% select(country_join, country_trase_id),
by = c("country" = "country_join")) %>%
left_join(df_eu, by = "country_trase_id") %>%
mutate(
country_trase_id = coalesce(country_trase_id, "UNKNOWN"),
economic_bloc_trase_id = coalesce(economic_bloc_trase_id, country_trase_id),
economic_bloc = ifelse(economic_bloc_trase_id == "EU", "EUROPEAN UNION", country))
message("Validate")
present_value_columns <- intersect(expected_value_columns, names(df))
stopifnot(all(sapply(df[present_value_columns], is.numeric)))
message(year, " rows=", nrow(df), " vol=", sum(df$vol, na.rm = TRUE))
message("fob=", sum(df$fob, na.rm = TRUE))
message("missing country_trase_id: ", sum(df$country_trase_id == "UNKNOWN", na.rm = TRUE))
message("missing country: ", sum(df$country == "UNKNOWN", na.rm = TRUE))
message("missing kabupaten_trase_id: ", sum(df$kabupaten_trase_id == "ID-XXXX", na.rm = TRUE))
message("missing province_trase_id: ", sum(df$province_trase_id == "ID-XX", na.rm = TRUE))
print(diagnose(df), n = Inf)
print(glimpse(df))
#final cleanup
df <- df %>%
select(-c(refinery_id)) %>%
mutate(
exporter_trase_id = ifelse(exporter_group == "DOMESTIC PROCESSING AND CONSUMPTION","", exporter_trase_id),
importer_group = ifelse(importer == "UNKNOWN", "UNKNOWN", importer_group),
branch = ifelse(branch == "UNKNOWN", "UNKNOWN", branch)
)
tmp_file <- tempfile(fileext = ".parquet")
write_parquet(df, tmp_file)
message("Write data")
put_object(
file = tmp_file,
object = output_key,
bucket = bucket,
region = region
)
message("Uploaded to s3://", bucket, "/", output_key)
S
Read country lookup
Read EU economic bloc lookup
process_year <- function(year, version = "v1.3") {
input_key <- paste0(
"indonesia/palm_oil/sei_pcs/", version, "/",
"results_", year, "_KM_0_MGD_30.csv"
)
output_key <- paste0(
"indonesia/palm_oil/sei_pcs/", version, "/",
"SEIPCS_INDONESIA_PALM_OIL_", year, "_KM_0_MGD_30.parquet"
)
message("Read data")
df <- s3read_using(
read_delim,
object = input_key,
bucket = bucket,
delim = ";",
show_col_types = FALSE
)
print(names(df))
message("Prepare results")
df <- df %>%
mutate(year = as.integer(year))
for (col in intersect(expected_value_columns, names(df))) {
df[[col]] <- as.numeric(df[[col]])
}
message("Join country ids and economic blocs")
df <- df %>%
select(-any_of(c("country_id", "country_trase_id", "economic_bloc_trase_id"))) %>%
mutate(country_join = toupper(trimws(country))) %>%
left_join(df_countries, by = "country_join") %>%
left_join(df_eu, by = "country_trase_id") %>%
mutate(
country_id = coalesce(country_id, "0"),
country_trase_id = coalesce(country_trase_id, "UNKNOWN"),
economic_bloc_trase_id = coalesce(economic_bloc_trase_id, country_trase_id)
) %>%
select(-country_join)
message("Validate")
present_value_columns <- intersect(expected_value_columns, names(df))
stopifnot(all(sapply(df[present_value_columns], is.numeric)))
message(year, " rows=", nrow(df), " vol=", sum(df$vol, na.rm = TRUE))
message("fob=", sum(df$fob, na.rm = TRUE))
message("missing country_trase_id: ", sum(df$country_trase_id == "UNKNOWN", na.rm = TRUE))
message("missing kabupaten_trase_id: ", sum(df$kabupaten_trase_id == "ID-XXXX", na.rm = TRUE))
message("missing province_trase_id: ", sum(df$province_trase_id == "ID-XX", na.rm = TRUE))
print(glimpse(df))
tmp_file <- tempfile(fileext = ".parquet")
write_parquet(df, tmp_file)
message("Write data")
put_object(
file = tmp_file,
object = output_key,
bucket = bucket
)
message("Uploaded to s3://", bucket, "/", output_key)
invisible(df)
}
#years <- c(2018, 2019, 2020, 2021, 2022, 2023)
years <- 2023
outs <- lapply(years, process_year)
Read data
[1] "branch" "certification" "commodity"
[4] "concession_trase_id" "country" "exporter_trase_id"
[7] "hs" "importer" "importer_group"
[10] "kabupaten_trase_id" "key" "mill_trase_id"
[13] "port_trase_id" "province_trase_id" "refinery_id"
[16] "refinery_trase_id" "year" "product_vol"
[19] "vol" "fob" "exporter"
[22] "exporter_group" "mill" "mill_group"
[25] "refinery_group" "refinery" "anonymize"
Prepare results
Error in FUN(X[[i]], ...): object 'expected_value_columns' not found
Traceback:
1. FUN(X[[i]], ...)
2. intersect(expected_value_columns, names(df))
3. .handleSimpleError(function (cnd)
. {
. watcher$capture_plot_and_output()
. cnd <- sanitize_call(cnd)
. watcher$push(cnd)
. switch(on_error, continue = invokeRestart("eval_continue"),
. stop = invokeRestart("eval_stop"), error = NULL)
. }, "object 'expected_value_columns' not found", base::quote(FUN(X[[i]],
. ...)))
df <- s3read_using(
read_delim,
object = input_key,
bucket = bucket,
delim = ";",
show_col_types = FALSE
)
#here can simplify if concession df include kabupaten
df <- df %>%
left_join(df_conc %>% select(ffb_code, kabupaten_trase_id, province_trase_id), by = "ffb_code") %>%
mutate(kabupaten_trase_id = if_else(is.na(kabupaten_trase_id), "ID-XXXX", kabupaten_trase_id),
province_trase_id = if_else(is.na(province_trase_id), "ID-XX", province_trase_id))
glimpse(df)
print("Compare result keys vs export keys")
names(df)
names(df_exports)
result_keys <- df %>%
count(commodity, exporter_trase_id, port_trase_id, name = "n_results")
export_keys <- df_exports %>%
count(commodity, exporter_trase_id, port_trase_id, name = "n_exports")
key_compare <- full_join(
result_keys,
export_keys,
by = c("commodity", "exporter_trase_id", "port_trase_id", "certification")
) %>%
mutate(
n_results = if_else(is.na(n_results), 0L, n_results),
n_exports = if_else(is.na(n_exports), 0L, n_exports)
)
print(key_compare %>% count(n_results > 0, n_exports > 0))
Error in `mutate()`:
ℹ In argument: `kabupaten_trase_id = if_else(is.na(kabupaten_trase_id),
"ID-XXXX", kabupaten_trase_id)`.
Caused by error:
! object 'kabupaten_trase_id' not found
Traceback:
1. mutate(., kabupaten_trase_id = if_else(is.na(kabupaten_trase_id),
. "ID-XXXX", kabupaten_trase_id), province_trase_id = if_else(is.na(province_trase_id),
. "ID-XX", province_trase_id))
2. mutate.data.frame(., kabupaten_trase_id = if_else(is.na(kabupaten_trase_id),
. "ID-XXXX", kabupaten_trase_id), province_trase_id = if_else(is.na(province_trase_id),
. "ID-XX", province_trase_id))
3. mutate_cols(.data, dplyr_quosures(...), by)
4. withCallingHandlers(for (i in seq_along(dots)) {
. poke_error_context(dots, i, mask = mask)
. context_poke("column", old_current_column)
. new_columns <- mutate_col(dots[[i]], data, mask, new_columns)
. }, error = dplyr_error_handler(dots = dots, mask = mask, bullets = mutate_bullets,
. error_call = error_call, error_class = "dplyr:::mutate_error"),
. warning = dplyr_warning_handler(state = warnings_state, mask = mask,
. error_call = error_call))
5. mutate_col(dots[[i]], data, mask, new_columns)
6. mask$eval_all_mutate(quo)
7. eval()
8. if_else(is.na(kabupaten_trase_id), "ID-XXXX", kabupaten_trase_id)
9. check_logical(condition)
10. is_logical(x)
11. .handleSimpleError(function (cnd)
. {
. local_error_context(dots, i = frame[[i_sym]], mask = mask)
. if (inherits(cnd, "dplyr:::internal_error")) {
. parent <- error_cnd(message = bullets(cnd))
. }
. else {
. parent <- cnd
. }
. message <- c(cnd_bullet_header(action), i = if (has_active_group_context(mask)) cnd_bullet_cur_group_label())
. abort(message, class = error_class, parent = parent, call = error_call)
. }, "object 'kabupaten_trase_id' not found", base::quote(NULL))
12. h(simpleError(msg, call))
13. abort(message, class = error_class, parent = parent, call = error_call)
14. signal_abort(cnd, .file)
15. signalCondition(cnd)
df <- s3read_using(
FUN = arrow::read_parquet,
object = "indonesia/palm_oil/sei_pcs/v1.2.2/SEIPCS_INDONESIA_PALM_OIL_2022_KM_0_MGD_30_PATCHED.parquet",
bucket = bucket
)
df <- as.data.frame(df)
glimpse(df)
Rows: 4,838,587
Columns: 35
$ anonymize [3m[90m<chr>[39m[23m "False"[90m, [39m"False"[90m, [39m"False"[90m, [39m"False"[90m, [39m"False"[90m, [39m"F…
$ branch [3m[90m<chr>[39m[23m "TRACEABILITY EXCESS LP 1"[90m, [39m"TRACEABILITY EXCES…
$ certification [3m[90m<chr>[39m[23m "ISCC"[90m, [39m"ISCC"[90m, [39m"ISCC"[90m, [39m"ISCC"[90m, [39m"ISCC"[90m, [39m"ISCC"[90m,[39m…
$ commodity [3m[90m<chr>[39m[23m "CPO"[90m, [39m"CPO"[90m, [39m"CPO"[90m, [39m"CPO"[90m, [39m"CPO"[90m, [39m"CPO"[90m, [39m"CPO"…
$ concession_trase_id [3m[90m<chr>[39m[23m "ID-PALM-CONCESSION-07493"[90m, [39m"ID-PALM-CONCESSION…
$ country [3m[90m<chr>[39m[23m "ITALY"[90m, [39m"ITALY"[90m, [39m"ITALY"[90m, [39m"ITALY"[90m, [39m"ITALY"[90m, [39m"I…
$ country_id [3m[90m<chr>[39m[23m "112"[90m, [39m"112"[90m, [39m"112"[90m, [39m"112"[90m, [39m"112"[90m, [39m"112"[90m, [39m"112"…
$ country_trase_id [3m[90m<chr>[39m[23m "IT"[90m, [39m"IT"[90m, [39m"IT"[90m, [39m"IT"[90m, [39m"IT"[90m, [39m"IT"[90m, [39m"IT"[90m, [39m"IT"[90m,[39m…
$ economic_bloc_trase_id [3m[90m<chr>[39m[23m "EU"[90m, [39m"EU"[90m, [39m"EU"[90m, [39m"EU"[90m, [39m"EU"[90m, [39m"EU"[90m, [39m"EU"[90m, [39m"EU"[90m,[39m…
$ exporter [3m[90m<chr>[39m[23m "PERKEBUNAN NUSANTARA III"[90m, [39m"PERKEBUNAN NUSANTA…
$ exporter_group [3m[90m<chr>[39m[23m "PTPN III"[90m, [39m"PTPN III"[90m, [39m"PTPN III"[90m, [39m"PTPN III"[90m,[39m…
$ exporter_group_id [3m[90m<chr>[39m[23m "2248132"[90m, [39m"2248132"[90m, [39m"2248132"[90m, [39m"2248132"[90m, [39m"22…
$ exporter_id [3m[90m<chr>[39m[23m "2248139"[90m, [39m"2248139"[90m, [39m"2248139"[90m, [39m"2248139"[90m, [39m"22…
$ exporter_trase_id [3m[90m<chr>[39m[23m "ID-TRADER-0274"[90m, [39m"ID-TRADER-0274"[90m, [39m"ID-TRADER-…
$ importer [3m[90m<chr>[39m[23m "CARGILL"[90m, [39m"CARGILL"[90m, [39m"CARGILL"[90m, [39m"CARGILL"[90m, [39m"CA…
$ importer_group [3m[90m<chr>[39m[23m "CARGILL"[90m, [39m"CARGILL"[90m, [39m"CARGILL"[90m, [39m"CARGILL"[90m, [39m"CA…
$ importer_group_id [3m[90m<chr>[39m[23m "1679384"[90m, [39m"1679384"[90m, [39m"1679384"[90m, [39m"1679384"[90m, [39m"16…
$ importer_id [3m[90m<chr>[39m[23m "1679385"[90m, [39m"1679385"[90m, [39m"1679385"[90m, [39m"1679385"[90m, [39m"16…
$ kabupaten_trase_id [3m[90m<chr>[39m[23m "ID-1203"[90m, [39m"ID-1203"[90m, [39m"ID-1208"[90m, [39m"ID-1207"[90m, [39m"ID…
$ key [3m[90m<chr>[39m[23m "EXPORTER > SMALL SET OF MILLS"[90m, [39m"EXPORTER > SM…
$ mill [3m[90m<chr>[39m[23m "PERKEBUNAN NUSANTARA III (HAPESONG)"[90m, [39m"PERKEBU…
$ mill_group [3m[90m<chr>[39m[23m "PTPN III"[90m, [39m"PTPN III"[90m, [39m"PTPN III"[90m, [39m"PTPN III"[90m,[39m…
$ mill_group_id [3m[90m<chr>[39m[23m "2248132"[90m, [39m"2248132"[90m, [39m"2248132"[90m, [39m"2248132"[90m, [39m"22…
$ mill_trase_id [3m[90m<chr>[39m[23m "ID-PALM-MILL-00570"[90m, [39m"ID-PALM-MILL-00570"[90m, [39m"ID…
$ port_trase_id [3m[90m<chr>[39m[23m "ID-PORT-0027"[90m, [39m"ID-PORT-0027"[90m, [39m"ID-PORT-0027"[90m,[39m…
$ province_trase_id [3m[90m<chr>[39m[23m "ID-12"[90m, [39m"ID-12"[90m, [39m"ID-12"[90m, [39m"ID-12"[90m, [39m"ID-12"[90m, [39m"I…
$ refinery [3m[90m<chr>[39m[23m "NOT REFINED"[90m, [39m"NOT REFINED"[90m, [39m"NOT REFINED"[90m, [39m"N…
$ refinery_group [3m[90m<chr>[39m[23m "NOT REFINED"[90m, [39m"NOT REFINED"[90m, [39m"NOT REFINED"[90m, [39m"N…
$ refinery_group_id [3m[90m<chr>[39m[23m "15518216"[90m, [39m"15518216"[90m, [39m"15518216"[90m, [39m"15518216"[90m,[39m…
$ refinery_id [3m[90m<chr>[39m[23m "15518215"[90m, [39m"15518215"[90m, [39m"15518215"[90m, [39m"15518215"[90m,[39m…
$ refinery_trase_id [3m[90m<chr>[39m[23m "ID-PALM-REFINERY-0000"[90m, [39m"ID-PALM-REFINERY-0000…
$ year [3m[90m<int>[39m[23m 2022[90m, [39m2022[90m, [39m2022[90m, [39m2022[90m, [39m2022[90m, [39m2022[90m, [39m2022[90m, [39m2022[90m,[39m…
$ fob [3m[90m<dbl>[39m[23m 6660.8982[90m, [39m12305.2124[90m, [39m49343.1267[90m, [39m38739.7701[90m, [39m…
$ product_vol [3m[90m<dbl>[39m[23m 7.5692025[90m, [39m13.9831959[90m, [39m56.0717349[90m, [39m44.0224660[90m, [39m…
$ vol [3m[90m<dbl>[39m[23m 7.5692025[90m, [39m13.9831959[90m, [39m56.0717349[90m, [39m44.0224660[90m, [39m…