--- title: "Example: Parallel Bulk API Calls" output: rmarkdown::html_vignette vignette: > %\VignetteIndexEntry{Example: Parallel Bulk API Calls} %\VignetteEngine{knitr::rmarkdown} %\VignetteEncoding{UTF-8} --- ```{r, include = FALSE} knitr::opts_chunk$set( collapse = TRUE, comment = "#>" ) ``` ## Overview Making hundreds or thousands of API calls is a common bottleneck in data pipelines. This example demonstrates parallelizing REST API calls to fetch data from multiple endpoints. **Use Case**: Data enrichment, web scraping, external data integration, geocoding **Computational Pattern**: I/O-bound parallel processing with rate limiting ## The Problem You need to enrich a dataset of 1,000 companies with external data from a REST API: - Company financial metrics - Stock prices - News sentiment - ESG scores Each API call takes 0.5-2 seconds due to network latency. Sequential execution would take 8-30 minutes. ## Setup ```{r setup, eval=FALSE} library(starburst) library(httr) library(jsonlite) ``` ## API Call Function Define a function that fetches data for one company: ```{r api-function, eval=FALSE} fetch_company_data <- function(ticker) { # Add small delay to respect rate limits Sys.sleep(runif(1, 0.1, 0.3)) # For demo purposes, we'll use a public API # In practice, replace with your actual API endpoint base_url <- "https://api.example.com/company" tryCatch({ # Fetch company info response <- httr::GET( paste0(base_url, "/", ticker), httr::timeout(10), httr::add_headers( "User-Agent" = "staRburst-example/1.0" ) ) # Check for success if (httr::status_code(response) == 200) { data <- httr::content(response, "parsed") # Extract relevant fields list( ticker = ticker, success = TRUE, company_name = data$name %||% NA, market_cap = data$marketCap %||% NA, pe_ratio = data$peRatio %||% NA, revenue = data$revenue %||% NA, employees = data$employees %||% NA, sector = data$sector %||% NA, timestamp = Sys.time() ) } else { # Handle API errors list( ticker = ticker, success = FALSE, error = paste("HTTP", httr::status_code(response)), timestamp = Sys.time() ) } }, error = function(e) { # Handle network errors list( ticker = ticker, success = FALSE, error = as.character(e), timestamp = Sys.time() ) }) } # Helper: null-coalescing operator `%||%` <- function(x, y) if (is.null(x)) y else x ``` ## Mock API for Demo Since we need a real API for testing, let's create a mock function that simulates API behavior: ```{r mock-api, eval=FALSE} # Mock function that simulates API with realistic delays fetch_company_data_mock <- function(ticker) { # Simulate network latency Sys.sleep(runif(1, 0.2, 0.8)) # Simulate occasional failures (5% rate) if (runif(1) < 0.05) { return(list( ticker = ticker, success = FALSE, error = "API timeout", timestamp = Sys.time() )) } # Generate mock data list( ticker = ticker, success = TRUE, company_name = paste("Company", ticker), market_cap = round(rnorm(1, 50e9, 20e9), 0), pe_ratio = round(rnorm(1, 25, 10), 2), revenue = round(rnorm(1, 10e9, 5e9), 0), employees = round(rnorm(1, 50000, 20000), 0), sector = sample(c("Technology", "Healthcare", "Finance", "Energy"), 1), timestamp = Sys.time() ) } ``` ## Generate Sample Data Create a list of 1,000 company tickers: ```{r sample-data, eval=FALSE} # Generate mock ticker symbols set.seed(123) n_companies <- 1000 tickers <- paste0( sample(LETTERS, n_companies, replace = TRUE), sample(LETTERS, n_companies, replace = TRUE), sample(LETTERS, n_companies, replace = TRUE) ) head(tickers) # [1] "NAL" "RPL" "OQM" "TYW" "AIT" "UMD" ``` ## Local Execution Run sequentially on local machine: ```{r local, eval=FALSE} cat(sprintf("Fetching data for %d companies locally...\n", length(tickers))) local_start <- Sys.time() local_results <- lapply(head(tickers, 50), fetch_company_data_mock) local_time <- as.numeric(difftime(Sys.time(), local_start, units = "secs")) cat(sprintf("[OK] Completed 50 calls in %.1f seconds\n", local_time)) cat(sprintf(" Estimated time for %d: %.1f minutes\n", n_companies, local_time * n_companies / 50 / 60)) ``` **Illustrative output**: ``` Fetching data for 1000 companies locally... [OK] Completed 50 calls in 24.3 seconds Estimated time for 1000: 8.1 minutes ``` ## Cloud Execution with staRburst Run all 1,000 API calls in parallel: Each fetch is short, so submitting 1,000 one-call tasks would be dominated by per-task S3 overhead (see the [Workload Shapes](https://starburst.ing/articles/workload-shapes.html) guide). Instead, batch the tickers into ~100 tasks and have each task fetch its chunk: ```{r cloud, eval=FALSE} cat(sprintf("Fetching data for %d companies on AWS...\n", n_companies)) # ~100 tasks, each fetching a chunk of tickers batches <- split(tickers, ceiling(seq_along(tickers) / ceiling(length(tickers) / 100))) results <- starburst_map( batches, function(batch) lapply(batch, fetch_company_data_mock), workers = 25, cpu = 1, memory = "2GB" ) results <- unlist(results, recursive = FALSE) # flatten to 1,000 results ``` **Illustrative output** (times/cost vary — see the Workload Shapes and Performance guides for measured numbers): ``` [Starting] Starting starburst cluster with 25 workers [Status] Processing 100 items with 25 workers [Starting] Submitting 100 tasks... [Wait] Progress: 100/100 [OK] Completed [Cost] Estimated cost: (printed per run) ``` ## Results Processing Analyze the fetched data: ```{r analysis, eval=FALSE} # Convert results to data frame results_df <- do.call(rbind, lapply(results, function(x) { if (x$success) { data.frame( ticker = x$ticker, company_name = x$company_name, market_cap = x$market_cap, pe_ratio = x$pe_ratio, revenue = x$revenue, employees = x$employees, sector = x$sector, success = TRUE, error = NA, stringsAsFactors = FALSE ) } else { data.frame( ticker = x$ticker, company_name = NA, market_cap = NA, pe_ratio = NA, revenue = NA, employees = NA, sector = NA, success = FALSE, error = x$error, stringsAsFactors = FALSE ) } })) # Summary success_rate <- mean(results_df$success) * 100 cat(sprintf("\n=== Results Summary ===\n")) cat(sprintf("Total companies: %d\n", nrow(results_df))) cat(sprintf("Successful fetches: %d (%.1f%%)\n", sum(results_df$success), success_rate)) cat(sprintf("Failed fetches: %d\n", sum(!results_df$success))) # Show sample of results cat("\n=== Sample Results ===\n") print(head(results_df[results_df$success, ], 10)) # Sector distribution cat("\n=== Sector Distribution ===\n") print(table(results_df$sector)) # Market cap summary cat("\n=== Market Cap Summary ===\n") cat(sprintf("Mean: $%.2fB\n", mean(results_df$market_cap, na.rm = TRUE) / 1e9)) cat(sprintf("Median: $%.2fB\n", median(results_df$market_cap, na.rm = TRUE) / 1e9)) cat(sprintf("Range: $%.2fB - $%.2fB\n", min(results_df$market_cap, na.rm = TRUE) / 1e9, max(results_df$market_cap, na.rm = TRUE) / 1e9)) ``` ## Performance For measured performance and when bursting is worth it, see `vignette("performance")` and `vignette("workload-shapes")`. ## Rate Limiting Considerations When working with real APIs: Each worker throttles itself **independently**, so the *aggregate* rate is the per-worker rate times the number of workers. To respect a global API limit, give each worker its share: `per_worker_rate = global_limit / workers`. ```{r rate-limiting, eval=FALSE} # per_worker_rate is calls/minute allowed to THIS worker fetch_with_rate_limit <- function(ticker, per_worker_rate) { delay <- 60 / per_worker_rate + runif(1, 0, 0.1) Sys.sleep(delay) fetch_company_data(ticker) } global_limit <- 100 # API allows 100 calls/minute total workers <- 10 per_worker <- global_limit / workers # 10 calls/min each results <- starburst_map( tickers, function(x) fetch_with_rate_limit(x, per_worker_rate = per_worker), workers = workers ) ``` **Rate limit calculation**: - API limit: 100 calls/minute (global) - Workers: 10 → each worker allowed 100 / 10 = 10 calls/minute - Aggregate throughput: 10 workers × 10 calls/min = **100 calls/minute** (at the limit) - Rule: `per_worker_rate = global_limit / workers` keeps the fleet under the cap ## Error Handling Best Practices ```{r error-handling, eval=FALSE} fetch_with_retry <- function(ticker, max_retries = 3) { for (attempt in 1:max_retries) { result <- fetch_company_data(ticker) if (result$success) { return(result) } # Exponential backoff if (attempt < max_retries) { Sys.sleep(2^attempt + runif(1, 0, 1)) } } # Return failure after all retries result$error <- paste("Failed after", max_retries, "retries:", result$error) return(result) } ``` ## When to Use This Pattern **Good fit**: - Many independent API calls (> 100) - Each call takes > 0.5 seconds - API allows concurrent requests - Transient failures are acceptable **Not ideal**: - Strict rate limits (< 10 calls/second total) - APIs that block concurrent requests - Very fast APIs (< 0.1 seconds per call) ## Running the Full Example The complete runnable script is available at: ```{r, eval=FALSE} system.file("examples/api-calls.R", package = "starburst") ``` Run it with: ```{r, eval=FALSE} source(system.file("examples/api-calls.R", package = "starburst")) ``` ## Next Steps - Replace mock function with your actual API - Implement proper authentication (API keys, OAuth) - Add request caching to avoid redundant calls - Monitor API usage and costs - Implement more sophisticated retry logic **Related examples**: - [Report Generation](example-reports.html) - Another I/O-bound parallel task - [Feature Engineering](example-feature-engineering.html) - Data enrichment patterns