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              <chr> "False", "False", "False", "False", "False", "F…
$ branch                 <chr> "TRACEABILITY EXCESS LP 1", "TRACEABILITY EXCES…
$ certification          <chr> "ISCC", "ISCC", "ISCC", "ISCC", "ISCC", "ISCC",…
$ commodity              <chr> "CPO", "CPO", "CPO", "CPO", "CPO", "CPO", "CPO"…
$ concession_trase_id    <chr> "ID-PALM-CONCESSION-07493", "ID-PALM-CONCESSION…
$ country                <chr> "ITALY", "ITALY", "ITALY", "ITALY", "ITALY", "I…
$ country_id             <chr> "112", "112", "112", "112", "112", "112", "112"…
$ country_trase_id       <chr> "IT", "IT", "IT", "IT", "IT", "IT", "IT", "IT",…
$ economic_bloc_trase_id <chr> "EU", "EU", "EU", "EU", "EU", "EU", "EU", "EU",…
$ exporter               <chr> "PERKEBUNAN NUSANTARA III", "PERKEBUNAN NUSANTA…
$ exporter_group         <chr> "PTPN III", "PTPN III", "PTPN III", "PTPN III",…
$ exporter_group_id      <chr> "2248132", "2248132", "2248132", "2248132", "22…
$ exporter_id            <chr> "2248139", "2248139", "2248139", "2248139", "22…
$ exporter_trase_id      <chr> "ID-TRADER-0274", "ID-TRADER-0274", "ID-TRADER-…
$ importer               <chr> "CARGILL", "CARGILL", "CARGILL", "CARGILL", "CA…
$ importer_group         <chr> "CARGILL", "CARGILL", "CARGILL", "CARGILL", "CA…
$ importer_group_id      <chr> "1679384", "1679384", "1679384", "1679384", "16…
$ importer_id            <chr> "1679385", "1679385", "1679385", "1679385", "16…
$ kabupaten_trase_id     <chr> "ID-1203", "ID-1203", "ID-1208", "ID-1207", "ID…
$ key                    <chr> "EXPORTER > SMALL SET OF MILLS", "EXPORTER > SM…
$ mill                   <chr> "PERKEBUNAN NUSANTARA III (HAPESONG)", "PERKEBU…
$ mill_group             <chr> "PTPN III", "PTPN III", "PTPN III", "PTPN III",…
$ mill_group_id          <chr> "2248132", "2248132", "2248132", "2248132", "22…
$ mill_trase_id          <chr> "ID-PALM-MILL-00570", "ID-PALM-MILL-00570", "ID…
$ port_trase_id          <chr> "ID-PORT-0027", "ID-PORT-0027", "ID-PORT-0027",…
$ province_trase_id      <chr> "ID-12", "ID-12", "ID-12", "ID-12", "ID-12", "I…
$ refinery               <chr> "NOT REFINED", "NOT REFINED", "NOT REFINED", "N…
$ refinery_group         <chr> "NOT REFINED", "NOT REFINED", "NOT REFINED", "N…
$ refinery_group_id      <chr> "15518216", "15518216", "15518216", "15518216",…
$ refinery_id            <chr> "15518215", "15518215", "15518215", "15518215",…
$ refinery_trase_id      <chr> "ID-PALM-REFINERY-0000", "ID-PALM-REFINERY-0000…
$ year                   <int> 2022, 2022, 2022, 2022, 2022, 2022, 2022, 2022,…
$ fob                    <dbl> 6660.8982, 12305.2124, 49343.1267, 38739.7701, …
$ product_vol            <dbl> 7.5692025, 13.9831959, 56.0717349, 44.0224660, …
$ vol                    <dbl> 7.5692025, 13.9831959, 56.0717349, 44.0224660, …