DZone
Thanks for visiting DZone today,
Edit Profile
  • Manage Email Subscriptions
  • How to Post to DZone
  • Article Submission Guidelines
Sign Out View Profile
  • Post an Article
  • Manage My Drafts
Newsletter
Log In / Join
Refcards Trend Reports
Events Video Library
Refcards
Trend Reports

Events

View Events Video Library

Performance

Performance refers to how well an application conducts itself compared to an expected level of service. Today's environments are increasingly complex and typically involve loosely coupled architectures, making it difficult to pinpoint bottlenecks in your system. Whatever your performance troubles, this Zone has you covered with everything from root cause analysis, application monitoring, and log management to anomaly detection, observability, and performance testing.

icon
Latest Premium Content
Trend Report
Observability and Performance
Observability and Performance
Refcard #290
Getting Started With Log Management
Getting Started With Log Management
Refcard #385
Observability Maturity Model
Observability Maturity Model

DZone's Featured Performance Resources

Why Time Series Databases Matter for Modern Enterprise Applications

Why Time Series Databases Matter for Modern Enterprise Applications

By Otavio Santana DZone Core CORE
Time is now a critical dimension in enterprise data. Modern applications generate a constant stream of events, such as payments, sensor readings, infrastructure metrics, user activity, logistics movements, market data, and application telemetry. The value exists not only in what happened, but also when it occurred, what came before it, and how the data changes over time. While traditional relational and NoSQL databases can store this information, increasing data volume and frequency make it essential to efficiently query recent states, historical windows, trends, and ordered events. Time-series databases fulfill these needs by treating time as a primary dimension. They make it easy to retrieve the latest measurements, examine historical intervals, aggregate data over time, and handle high-frequency writes. For enterprise applications, time-series databases now support not only monitoring and IoT, but also observability, financial systems, industrial platforms, logistics, energy, and any architecture where tracking data evolution is as important as the data itself. Why Time Series Databases Matter At first glance, time-series data appears similar to data stored in relational or NoSQL databases. SQL tables can include timestamp columns, document databases can store dated events, and key-value stores can use time-based keys. The distinction appears when time becomes the primary access pattern rather than a secondary attribute. If typical queries include “what is the latest value?”, “what happened during this interval?”, or “how did this metric change over time?”, a time-series model is often a better fit. Relational databases can support these queries, but they frequently require complex indexes, partitions, retention policies, and aggregation logic as data volume increases. Document and key-value databases typically require the application to manage temporal structure. Time-series databases are designed for append-heavy workloads, ordered data, time-window queries, downsampling, retention, and time-based aggregation. They do not replace SQL or general-purpose NoSQL databases, but serve as a specialized solution when tracking data changes over time is essential. A useful way to think about the difference is this: Plain Text Relational database: "What is the current state of this entity?" Document database: "What does this aggregate/document look like?" Key-Value database: "What value is associated with this key?" Time Series database: "How has this value changed over time, and what is happening now?" Consider a temperature sensor. In a relational model, we could create a table like: SQL CREATE TABLE sensor_reading ( timestamp TIMESTAMP, sensor VARCHAR(100), temperature DOUBLE ); This approach is functional, but typical workloads often require queries such as: Plain Text Give me the latest value for sensor-01. Give me the last 100 measurements. Give me the average temperature every five minutes. Show me the period where the temperature exceeded 30°C. Compare today's measurements with yesterday's. Delete or compress measurements older than six months. These are not isolated queries; they represent the typical data lifecycle. Time-series databases are specifically designed to support this pattern. Beyond IoT IoT is a clear example, as sensors naturally generate timestamped measurements. However, this model is common across many enterprise systems. Observability and infrastructure monitoring are classic use cases. Measurements such as CPU usage, memory consumption, request latency, error rates, queue depth, database connections, and network throughput are sampled repeatedly over time. Rather than focusing on a single point, teams typically seek trends, spikes, rolling averages, anomalies, and interval comparisons. Financial systems also generate highly temporal data. Market prices, exchange rates, trades, account activity, risk measurements, and portfolio valuations are all chronological streams. Trading applications may require the latest price, recent ticks, or price ranges within specific windows. Banking systems regularly retain transactional data in relational databases while sending temporal metrics and activity streams to time-series stores for analysis. Logistics and transportation also benefit from this approach. Vehicle position, speed, fuel consumption, delivery status, warehouse temperature, and route telemetry are all continuously changing metrics. Relevant queries include: Plain Text Where was the vehicle during the last hour? How has fuel consumption changed during this route? When did the delivery temperature leave the acceptable range? What was the most recent location reported by each vehicle? Energy and utilities are also inherently time-oriented. Smart meters, electricity consumption, voltage, water usage, solar generation, battery charge, and grid load all generate continuous measurements. While billing may remain in a relational system, raw measurements and their aggregation suit time-series databases well. Application and business metrics further demonstrate that time-series data goes beyond infrastructure. Enterprise applications may record: orders per minutepayments approved per minutefailed logins per houractive userscheckout latencyinventory changesfraud scoresmessages processedAPI calls by customer These are business signals, not simply technical telemetry. Persisting them over time enables organizations to understand both system and business processes as they evolve. Event-driven and distributed architectures present another important use case. Microservices continuously emit events and measurements. While event stores and message brokers are effective for transporting and preserving events, they are not optimized for queries such as: Plain Text What was the average latency for this service during the last 30 minutes? What is the latest measurement from every region? How has throughput changed since the deployment? Which customers experienced the highest error rate over the last hour? A time-series database can complement Kafka, Pulsar, or other event backbones by providing a queryable temporal view of these streams. This distinction is important: selecting a time-series database does not require replacing relational databases, document stores, or event brokers. Modern enterprise architectures are increasingly polyglot. Relational databases may remain the system of record, document databases may manage flexible aggregates, key-value databases may provide fast lookups, and time-series databases may handle continuously changing measurements. The architectural advantage lies in choosing a data model that corresponds to the application's key questions. Java and Time Series Traditionally, using time-series databases in Java required working with database-specific drivers and APIs. Developers needed to learn unique programming models, configuration styles, query APIs, and object-mapping strategies for each database. This increased implementation complexity and made maintenance challenging, especially when multiple time-series technologies were used. Eclipse JNoSQL 1.1.18 addresses this complexity by introducing Time Series as a supported database type. This release adds support for InfluxDB 3, Apache IoTDB, and QuestDB, along with a new Jakarta NoSQL Template specialization: TimeSeriesTemplate. Java @Inject private TimeSeriesTemplate template; The mapping model remains consistent. Time-series entities continue to use Jakarta NoSQL annotations: Java @Entity public class SensorReading { @Id private Instant timestamp; @Column private String sensor; @Column private double value; // constructors, getters, and setters } As with other Jakarta NoSQL database types, the primary difference lies in the data model. In time-series models, the identifier usually represents the temporal dimension. In this example, @Id is an Instant marking when the measurement occurred. The programming model remains familiar: Java @Inject private TimeSeriesTemplate template; SensorReading reading = new SensorReading( Instant.now(), "sensor-1", 21.5); template.insert(reading); Optional<SensorReading> result = template.find(SensorReading.class, reading.getTimestamp()); List<SensorReading> readings = template.select(SensorReading.class) .where("sensor").eq("sensor-1") .result(); The same model applies to Jakarta Data repositories: Java @Repository public interface SensorReadingRepository extends BasicRepository<SensorReading, Instant> { List<SensorReading> findBySensor(String sensor); } Then the repository can be injected normally: Java @Inject private SensorReadingRepository repository; The key improvement is not Java’s ability to connect to time-series databases, as native drivers have always enabled this. The real value is that developers can now use a consistent Jakarta NoSQL and Jakarta Data programming model across different time-series databases, reducing database-specific code within applications. Conclusion Time-series databases matter because many modern systems are no longer only about storing the current state of data, but about understanding how that data changes over time. From observability and finance to logistics, energy, IoT, and business metrics, a time-oriented model can simplify both the architecture and the queries when temporal behavior is central to the problem. More
Part 3: End-to-End Tracing and Observability Across Goose, agentgateway, and Quarkus

Part 3: End-to-End Tracing and Observability Across Goose, agentgateway, and Quarkus

By Daniel Oh DZone Core CORE
Enterprise context — Acme FinServ. SOC 2 CC7 (system monitoring) requires that Acme can detect and investigate anomalous activity. When an agent-driven workflow touches customer data at 2 AM, "we have logs somewhere" is not an answer an auditor accepts. The distributed trace built in this part is the forensic evidence trail: a single trace ID that ties the Goose prompt to every agentgateway policy decision and every Quarkus tool call, so a post-incident review can reconstruct exactly which agent did what, in what order, and how long each governed hop took. The Core Problem In Part 1, we built a Quarkus MCP tool server. In Part 2, we secured it with agentgateway's JWT authentication, RBAC, and ExtMCP guardrails. The architecture works — but when something goes wrong in production, you're flying blind. Agentic workflows are fundamentally different from traditional request-response APIs. A single user prompt like "Debug customer CUST-4091" triggers a multi-round-trip loop: Goose calls tools/list to discover available toolsThe LLM selects getCustomerStatus and Goose sends tools/callThe LLM reads the response, sees primaryRegion: US-EAST-1, and chains a second tools/call to getZoneHealthLogsThe LLM correlates both results and generates a diagnostic summary Each of these hops crosses process boundaries: Goose → agentgateway → Quarkus. Without distributed tracing, you see four isolated HTTP requests in your access logs. You cannot tell they belong to the same agentic workflow. When step 3 takes 12 seconds instead of 200ms, you have no waterfall to pinpoint whether the latency came from agentgateway policy evaluation, Quarkus bean validation, or a slow downstream call. This creates black holes in telemetry dashboards — the exact gap that autonomous agents exploit to degrade silently. The Solution: W3C Trace Context Across All Three Layers The fix is standard distributed tracing, applied to the MCP transport layer: agentgateway exports spans for every proxied MCP request and propagates traceparent headers to the backend.Quarkus with quarkus-opentelemetry picks up the incoming traceparent, creates child spans for tool execution and bean validation, and exports them to the same Jaeger instance.Jaeger correlates both sides into a single trace waterfall — one view from agent prompt to tool result. Prerequisites Everything from Parts 1 and 2, plus: Podman – for running Jaeger (podman compose) Verify Podman is available: Shell podman --version Step 1: Launching the Observability Backend We use Jaeger v2 as both the OTLP collector and the trace UI. A single container accepts traces from agentgateway on port 4317 (OTLP gRPC) and from Quarkus on port 4318 (OTLP HTTP), and serves the query UI on port 16686. Shell cd part3-observability podman compose up -d This starts Jaeger v2 with OTLP collection enabled by default. Verify it's running: Shell curl -sf http://localhost:16686/ > /dev/null && echo "Jaeger UI is ready" Open http://localhost:16686 — you'll see an empty Jaeger UI. We'll populate it with MCP traces in the following steps. Production Alternative: Grafana Tempo For production deployments, replace Jaeger with Grafana Tempo backed by object storage (S3/GCS). The OTLP endpoint stays the same — only the compose.yml changes. Grafana provides richer dashboards, alerting, and long-term trace retention. Step 2: Enabling OpenTelemetry in Quarkus Add the quarkus-opentelemetry extension to Part 1's pom.xml: Properties files <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-opentelemetry</artifactId> </dependency> Configure the exporter in application.properties: Properties files # OpenTelemetry quarkus.otel.service.name=customer-tools quarkus.otel.exporter.otlp.traces.endpoint=http://localhost:4318 quarkus.otel.exporter.otlp.traces.protocol=http/protobuf quarkus.otel.traces.sampler=always_on quarkus.otel.traces.suppress-non-application-uris=false PropertyPurposeservice.nameIdentifies this service in Jaeger's service dropdowntraces.endpointOTLP HTTP receiver — Jaeger's port 4318 (base URL only; Quarkus appends /v1/traces)traces.protocolhttp/protobuf — Quarkus uses its Vert.x-based HTTP exportertraces.sampleralways_on — sample every span (reduce in production)suppress-non-application-urisfalse — include MCP endpoint spans (they'd be filtered otherwise) When no OTLP collector is running (Parts 1 and 2 without Jaeger), Quarkus logs a connection warning, but the MCP server works normally. When the collector IS running (Part 3), traces flow automatically. Zero code changes to the MCP tools. Rebuild Part 1: Shell cd part1-quarkus-mcp mvn package -DskipTests What Quarkus Auto-Instruments With quarkus-opentelemetry on the classpath and the SDK enabled, Quarkus automatically creates spans for: LayerSpan NameWhat It CapturesHTTP serverPOST /mcpInbound MCP request with method, status, latencyCDI beansCustomerServiceTools.getCustomerStatusTool execution time within the MCP handlerBean ValidationHibernateValidatorParameter validation before tool logic runsREST clientOutbound HTTP callsAny downstream API calls (future extensions) No @WithSpan annotations needed. The Quarkus OpenTelemetry extension instruments the reactive pipeline automatically. Step 3: Configuring W3C Trace Context in agentgateway agentgateway supports native OpenTelemetry trace export. Add a tracing block to the gateway configuration: Properties files config: adminAddr: localhost:15000 tracing: otlpEndpoint: http://localhost:4317 otlpProtocol: grpc randomSampling: 1.0 FieldPurposeotlpEndpointOTLP receiver — Jaeger's port 4317otlpProtocolgrpc for OTLP/gRPC (also supports http)randomSamplingSample 100% of traces (reduce to 0.01–0.1 in production) How Trace Propagation Works When agentgateway receives an MCP request: Creates a root span for the proxy operation (e.g., agentgateway.mcp.proxy)Injects a traceparent header into the forwarded request to Quarkus:traceparent: 00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01Quarkus reads the traceparent, creates a child span under the same trace ID, and records tool executionBoth spans export to Jaeger via OTLP, where they appear as a single correlated trace This is standard W3C Trace Context propagation — the same mechanism used across all OpenTelemetry-instrumented services. Configuration Files Part 3 provides two agentgateway configurations: ConfigUse Caseconfig-traced.yamlTracing only — proxy + OTLP export, no security layersconfig-traced-guardrails.yamlTracing + ExtMCP guardrails — observe the guardrail evaluation spans too Step 4: Running the Interactive Demo Start all services with the one-command script: Shell cd part3-observability ./start-all.sh The script starts Jaeger, Quarkus (with OTel enabled), and agentgateway (with trace export), then launches the demo SPA on :8890. Open the MCP Observability Console at http://localhost:8890/index.html and walk through the three demo steps: Initialize – Establishes an MCP session through agentgateway. The architecture diagram animates the trace propagation: root span creation in agentgateway, traceparent injection, child span in Quarkus, and OTLP export to Jaeger.List Tools – Discovers all 5 tools through the traced proxy. The trace waterfall panel shows the agentgateway proxy span and the Quarkus HTTP span side by side with timing.Multi-Tool Workflow – Simulates Goose's multi-turn reasoning: getCustomerStatus (finds region US-EAST-1) → getZoneHealthLogs (checks zone health) → getSLACompliance (correlates SLA metrics). Each step generates a full trace with waterfall visualization. The stat tiles track traces generated, spans collected, and Jaeger status. Click Open Jaeger to view the real trace waterfalls in the Jaeger UI at http://localhost:16686. Step 5: Generating Traces via CLI To generate additional traces manually, simulate a multi-turn agentic workflow: Shell # Step 1: Initialize MCP session export MCP_SESSION_ID=$(curl -s -D - http://localhost:3000/mcp \ -H "Content-Type: application/json" \ -H "Accept: application/json, text/event-stream" \ -H "MCP-Protocol-Version: 2025-03-26" \ -d '{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-03-26","capabilities":{},"clientInfo":{"name":"curl","version":"1.0"}}' \ | grep -i "mcp-session-id:" | sed 's/.*: //' | tr -d '\r') # Step 2: Discover tools curl -s http://localhost:3000/mcp \ -H "Content-Type: application/json" \ -H "Accept: application/json, text/event-stream" \ -H "MCP-Protocol-Version: 2025-03-26" \ -H "mcp-session-id: $MCP_SESSION_ID" \ -d '{"jsonrpc":"2.0","id":2,"method":"tools/list","params":{}' \ | grep '^data: ' | sed 's/^data: //' | jq . # Step 3: Agent calls getCustomerStatus (first tool invocation) curl -s http://localhost:3000/mcp \ -H "Content-Type: application/json" \ -H "Accept: application/json, text/event-stream" \ -H "MCP-Protocol-Version: 2025-03-26" \ -H "mcp-session-id: $MCP_SESSION_ID" \ -d '{"jsonrpc":"2.0","id":3,"method":"tools/call","params":{"name":"getCustomerStatus","arguments":{"customerId":"CUST-4091"}}' \ | grep '^data: ' | sed 's/^data: //' | jq . # Step 4: Agent chains getZoneHealthLogs based on the region from step 3 curl -s http://localhost:3000/mcp \ -H "Content-Type: application/json" \ -H "Accept: application/json, text/event-stream" \ -H "MCP-Protocol-Version: 2025-03-26" \ -H "mcp-session-id: $MCP_SESSION_ID" \ -d '{"jsonrpc":"2.0","id":4,"method":"tools/call","params":{"name":"getZoneHealthLogs","arguments":{"zoneId":"US-EAST-1"}}' \ | grep '^data: ' | sed 's/^data: //' | jq . # Step 5: Agent fetches SLA compliance for correlation curl -s http://localhost:3000/mcp \ -H "Content-Type: application/json" \ -H "Accept: application/json, text/event-stream" \ -H "MCP-Protocol-Version: 2025-03-26" \ -H "mcp-session-id: $MCP_SESSION_ID" \ -d '{"jsonrpc":"2.0","id":5,"method":"tools/call","params":{"name":"getSLACompliance","arguments":{"serviceId":"api-gateway"}}' \ | grep '^data: ' | sed 's/^data: //' | jq . Each of these requests generates a trace that flows through agentgateway into Quarkus and lands in Jaeger. Step 6: Visualizing the Trace Waterfall in Jaeger Open http://localhost:16686 in your browser. Finding Traces In the Service dropdown, select customer-tools (Quarkus) or agentgatewayClick Find TracesClick on any trace to open the waterfall view Reading the Waterfall A typical tools/call trace shows the following span hierarchy: Shell agentgateway.mcp.proxy [12ms] └─ POST /mcp [8ms] ← Quarkus HTTP server └─ CustomerServiceTools.getCustomerStatus [2ms] ← CDI tool execution SpanServiceWhat It Tells Youagentgateway.mcp.proxyagentgatewayTotal proxy overhead including policy evaluationPOST /mcpcustomer-toolsQuarkus HTTP handling time for the MCP requestgetCustomerStatuscustomer-toolsPure tool execution time (business logic) What to Look For Proxy overhead: The gap between the agentgateway span and the Quarkus span shows network + policy evaluation time. If this grows, check guardrail server latency.Validation time: Bean Validation spans appear before tool execution. Regex-heavy patterns like ^CUST-[0-9]{4,8}$ are fast, but complex validators on large payloads can add latency.Multi-turn correlation: When Goose chains multiple tool calls (e.g., getCustomerStatus → getZoneHealthLogs), each appears as a separate trace. The mcp-session-id tag lets you filter all traces belonging to one agent session.Error traces: Failed validations (invalid customer ID format) or guardrail rejections (blocked poison payloads) produce error spans with exception details. Connecting Goose for Real Traces Launch Goose pointed at agentgateway and prompt a multi-tool workflow: Shell goose session "Debug customer CUST-4091 — check their account status, then pull health logs for their region and SLA compliance for api-gateway." This generates a burst of correlated traces in Jaeger showing Goose's multi-turn tool orchestration from the proxy layer down to individual tool execution spans. What We Achieved Starting from the secured architecture in Part 2, we added full observability without changing any MCP tool code: LayerWhat We AddedConfig ChangeQuarkusquarkus-opentelemetry dependencypom.xml + application.propertiesagentgatewaytracing block in config YAMLconfig-traced.yamlObservability backendJaeger all-in-one via Podman Composecompose.yml The entire stack runs locally with a single ./start-all.sh command and produces end-to-end trace waterfalls in Jaeger. Production Considerations ConcernLocal (this tutorial)ProductionTrace backendJaeger all-in-one (in-memory)Grafana Tempo + object storageSampling rate100% (default: 1.0)1-10% or adaptive samplingTrace retentionContainer lifetimeDays/weeks in durable storageAlertingManual Jaeger inspectionGrafana alerting on span latency/error rateMetricsTraces onlyAdd Prometheus + quarkus-micrometer for RED metrics Coming Up in Part 4 With tracing in place, you can now see every MCP tool call flowing through the system. In Part 4, we will move beyond single-agent tool calls to multi-agent orchestration — using the Agent-to-Agent (A2A) protocol to coordinate autonomous agents that can delegate work, enforce governance via AGENTS.md, and call back into our MCP tool services. More
Gossips on Cryptography: Part 4
Gossips on Cryptography: Part 4
By Sahil Aggarwal
The Telemetry Tax: Architecting Zero-Allocation Event Observability at 15B+ Daily Event Scale
The Telemetry Tax: Architecting Zero-Allocation Event Observability at 15B+ Daily Event Scale
By Brindal Patel
The Silent Container Death: A TCP Dial That Never Times Out
The Silent Container Death: A TCP Dial That Never Times Out
By Alexander Fo
Predict, Repeat, Improve: Deterministic Simulation Testing Explained
Predict, Repeat, Improve: Deterministic Simulation Testing Explained

It’s 2 AM. Your phone buzzes, the on-call alert flashes, and suddenly you are staring at a production outage that makes no sense. Following the logs, you get a hint: when you rerun the same scenario in staging, everything behaves perfectly. None of the quality gates/QA pipelines catch it, chaos experiments didn’t reproduce it, and now a ghost is chased that only appears when the system is under real-world pressure. Distributed systems are notorious for these “phantom failures” — rare timing-dependent bugs that surface unpredictably and vanish just as quickly. They are the kind of dreaded incidents that keep engineers awake at night because they are unreproducible. Take a real-world example: A service once crashed because two nodes tried to become leader at the exact same millisecond. In staging, the timing never aligned that way, so the bug remained invisible. But in production, under heavy load, it just happens, sending the system into chaos. Engineers spent days trying to recreate the failure, but without a deterministic replay, it's pure luck to get a reliable reproduction. Even when it happens, engineers may not be sure what caused it or how to reproduce it deterministically. Enter deterministic simulation testing (DST). DST builds a fully controlled, re-playable simulation of your system’s world — nodes, clients, clocks, network delays, partitions — so that even the most elusive bugs can be identified, captured, replayed, and studied. In this article, we will uncover how DST can transform those unpredictable 2 AM incidents into predictable, debuggable coordinates — giving you a new way to tame the chaos of distributed systems. Deterministic Simulation Testing Definition Deterministic simulation testing (DST) is a software testing methodology that places the system under test within a fully controlled, simulated environment. All sources of non-determinism — system clock, thread scheduling, network, disk — are intercepted and made deterministic. The key property is that, for a given initial seed and configuration, the entire execution is reproducible. The same sequence of events, faults, and outcomes will occur on every run with that seed. Key Concepts Breaking down the above definition, below are the key concepts for DST: Determinism → The system’s behavior is purely dependent upon its initial state and the seed. All non-deterministic sources are simulated to achieve determinism.Simulation → The system is run in a virtual environment that can simulate faults and control the passage of time. E.g., controlled clock skew introduced across various nodes, added network delays to achieve out-of-order event delivery.Reproducibility → Any failure or bug found during simulation can be reliably reproduced by rerunning the simulation with the same seed.Scenario Exploration → By varying the seed and/or simulation parameters, DST systematically explores a vast range of possible execution paths and failure scenarios. How DST Works Let's consider a simple scenario where two users update and read the same record in a very short interval. User 1 updates a record with a new value at time instance T0, and User 2 reads the same record at time instance T1. Note that the interval between T0 & T1 stays the same. In an ideal case (i.e., scenario 1), the new value is updated or written immediately, i.e., without any delay. Thus, when User 2 reads the same record at T1, it is able to read the latest value. In scenario 2 suppose the write is delayed due to network partitioning, disk write etc. User 2 thus sees the old value of the record even if it reads the value at the same time instance T1. Although a stale read may look trivial, it may lead to workflow halt, process crash, etc. in a complex real-world system. Imagine what could happen in a real-world system where: Multiple processes are scheduled for execution, within and across nodes.Multiple network calls are made between several nodes.Multiple operations are performed by several distributed processes on a single disk. Traditional testing strategies or frameworks are inherently constrained and thus can’t simulate such delays or faults. Because of this, it's nearly impossible to identify, catch, reproduce, or debug issues arising from such situations — rendering them unreliable or, at best, non-deterministic. To achieve determinism, the testing framework must take total control over the environment to intercept and manage all external interactions as described below: Controlled scheduling → Instead of relying on the operating system’s unpredictable thread/coroutine scheduler, the simulator provides its own deterministic scheduler. Thus ensuring various scheduling combinations are simulated.I/O mocking → All network calls, disk writes, and clock queries are routed through the simulator, allowing it to inject latency, drop packets, or change the time (e.g., clock skew).Single-threaded execution → Many DST frameworks run the entire distributed system stack within a single thread, completely stripping away the chaotic, unrepeatable nature of multi-threading. Thus, by eliminating real-world “flakiness,” DST allows developers to reproduce chaotic distributed system bugs with perfect precision, thanks to its inherent ability to replay any failing execution: Seed-based replay → The same seed reproduces the exact sequence of events, making debugging tractable.Time-travel debugging → Some platforms (e.g., Flashback) allow stepping backward and forward through execution, inspecting state at any point for a granular view of the system. DST Implementation Approaches and Architecture Patterns Below are two approaches for DST. Pluggable Non-Determinism Design the system so that all non-deterministic components (clocks, I/O, etc.) are pluggable. This strategy is used by TigerBeetle. This requires: Abstracting all system interactions behind interfaces.Providing both real and simulated implementations.Ensuring that the same codebase can run in both production and simulation by swapping implementations at startup. Pros: Deep control and minimal divergence between test and production code. Suitable for greenfield systems. Cons: Requires significant upfront design and is challenging to retrofit into existing systems. Deterministic Hypervisors and Emulation A more recent and flexible approach is to run unmodified binaries inside a deterministic hypervisor or emulation layer. This strategy is used by Hermit and Weave. Pros: Can test existing systems without code changes; language-agnostic; simulates the entire stack. Cons: May have performance overhead; some system behaviors may escape determinism if not fully intercepted. Benefits System employing DST benefits as below: Identify and reproduce rare failures → DST allows engineers to replay the exact sequence of events that led to a bug. This eliminates the frustration of “flaky” issues that appear inconsistently, making debugging far more reliable. Moreover, DST helps find bugs in execution paths unreachable by example-based tests.Improved developer productivity → Bugs are easier to reproduce, debug, and fix; less time spent on war rooms and emergency triage.Improved confidence in correctness → DST validates critical invariants (like consensus, failover, or transaction consistency) under controlled simulations. Engineers gain assurance that core distributed protocols behave as expected even under stress. Thus, preventing rare bugs from reaching production, increasing system uptime and user trust.Scalable debugging for complex systems → In microservice or event-driven architectures, DST helps tame the exponential growth of possible interleavings by focusing on deterministic seeds. This makes large-scale debugging more tractable. Challenges and Limitations While DST provides unparalleled confidence, it requires significant architectural investment. Retrofitting DST into existing systems may require significant refactoring. It can be highly intrusive, requiring developers to write custom code or frameworks, as production code often cannot rely on external third-party libraries that invoke un-mocked I/O or system calls. Moreover, DST requires careful modeling of external systems to avoid missing integration bugs. Ensuring sufficient coverage without combinatorial explosion is a major challenge — especially in modern systems with multiple integration points. Below are gaps in tooling that limit DST outcomes: Language and platform support → Not all languages and runtimes have mature DST frameworks.Hypervisor limitations → Deterministic hypervisors may not support all system calls or hardware features. DST Comparison and Applicability DST vs. Chaos Engineering DST is proactive and enables perfect reproducibility. It is best suited for development and pre-production, catching bugs before they reach users. It can simulate production chaos in minutes, and every failure is a permanent regression. Chaos engineering is reactive, non-deterministic, and validates the behavior of real deployments. It is essential for catching issues arising from real infrastructure, misconfigurations, or dependencies that simulation cannot model. However, it cannot guarantee coverage or reproducibility, and carries the risk of impacting users. In essence, both DST and chaos engineering are complementary to each other and are necessary for comprehensive reliability. DST in Functional vs. Performance Testing Functional Testing DST is ideally suited for functional testing. Validates correctness under all possible interleavings, failures, and workloads.Checks invariants, safety properties, and liveness under stress.Finds rare, timing-dependent bugs that are invisible to example-based tests. Performance Testing DST is not primarily designed for performance testing. The simulated environment does not reflect real hardware performance, network latency, or throughput.Time is virtualized and compressed; I/O is in-memory.Performance metrics (latency, throughput) measured in simulation may not correspond to real-world values. However, DST can be used to: Validate performance-related invariants (e.g., absence of deadlocks, progress under load).Simulate pathological scenarios (e.g., extreme contention, resource exhaustion) to observe system behavior. Recommendation Combine DST for functional correctness with real-world performance and benchmarking suites for comprehensive validation. DST Applicability to AI/ML and Agentic Systems AI/ML systems, especially those based on large language models (LLMs) and agentic workflows, are fundamentally non-deterministic. This makes traditional testing and debugging extremely difficult, with “heisenbugs” that vanish when observed. DST can be adapted to AI/ML systems by creating controlled, simulated environments for agents to operate in. Or using a hybrid approach of combining deterministic components (rule-based logic) with LLM-driven reasoning, using seeds to replay failures. Case Studies FoundationDB, with its deterministic simulator tool, achieved legendary reliability by running trillions of simulated CPU-hours, finding and fixing every known bug before production.TigerBeetle built a Viewstamped Operation Replication simulator (VOPR) to simulate financial transaction systems, catching subtle bugs in consensus and replication.Ethereum Merge used Antithesis to test the transition to Proof-of-Stake, simulating multiple client implementations in a deterministic environment. Conclusion Deterministic simulation testing (DST) represents a paradigm shift in the testing and validation of distributed systems. By enabling exhaustive, reproducible exploration of the vast state space of concurrent, failure-prone systems, DST empowers engineers to find and fix the rarest and most pernicious bugs before they reach production. Its integration with property-based testing, fuzzing, and fault injection, combined with advances in deterministic hypervisors and simulation frameworks, has made DST accessible to a growing range of systems and organizations. While DST requires significant engineering investment, careful system design, and ongoing maintenance, its benefits in reliability, developer productivity, and user trust are profound. As distributed systems continue to grow in complexity and AI/ML systems become more agentic and autonomous, the need for rigorous, deterministic validation will only intensify. The future of DST lies in deeper integration with formal methods, smarter state-space exploration, and broader applicability to AI/ML and hybrid systems. Organizations that embrace DST, alongside complementary techniques like chaos engineering and formal verification, will be best positioned to deliver robust, trustworthy, and resilient distributed systems in the years ahead. DST Tools and Frameworks Deterministic simulation testing (DST) tooling is still a niche but growing ecosystem. Each has a unique focus — ranging from language-level deterministic runtimes to full-stack hypervisor-based reproducibility. Based on the specific needs a single or combination of them can be picked up. References and Further Reads Taming Chaos — DSTSquashing the Heisenbug with DSTAntithesis — DSTPhil Eaton — DSTRedstone — DST FrameworkResonate — DSTJespen | TickLoom

By Ammar Husain DZone Core CORE
Jakarta Batch in Practice: Reliable Chunk-Oriented Processing for Enterprise Workloads
Jakarta Batch in Practice: Reliable Chunk-Oriented Processing for Enterprise Workloads

Batch processing remains vital because many business operations aren't suited to interactive requests. Tasks such as recalculating prices, reconciling transactions, migrating records, generating reports, processing invoices, reclassifying customers, or applying rules across millions of records may require considerable time. Handling these as standard requests leads to fragile systems, increased user wait times, frequent timeouts, challenging retries, and possible data inconsistencies. A batch model handles large workloads predictably, incrementally, and with control over progress and recovery. Rather than processing a massive operation as a single loop, batch processing uses jobs, steps, chunks, checkpoints, filtering, and restartability. This approach separates long-running data tasks from the user experience while delivering a structured execution model. In this article, we will focus on Jakarta Batch and examine its sustained relevance for modern enterprise applications. Why Batch Processing Still Matters in Enterprise Systems Modern applications offer various methods for background processing, such as message queues, event-driven architectures, schedulers, reactive pipelines, and distributed stream-processing platforms. While each addresses specific needs, batch processing is most effective when operations have a defined start and end, involve a known or discoverable dataset, and require controlled execution, progress tracking, restartability, or periodic processing. Batch processing remains essential in enterprise systems. Workloads such as financial reconciliation, billing, payroll, reporting, data migration, regulatory processing, catalog updates, and large-scale reclassification are still prevalent. In these scenarios, the priority is to process large volumes of work safely and predictably, rather than responding to individual events quickly. Batch provides a model specifically designed for these requirements. How Jakarta Batch Works Jakarta Batch organizes background processing into jobs and steps. A job defines the overall batch operation, while each step represents a specific stage. In chunk-oriented processing, a step follows a simple pipeline: read, process, write, and repeat until it processes all input. The Jakarta Batch runtime manages this lifecycle so application code can focus on reading, transforming, and persisting data. A job is the top-level unit of execution and represents a complete business operation, such as importing records, recalculating customer classifications, processing invoices, or reconciling transactions. Jobs can accept parameters at startup, allowing the same batch definition to run with different inputs or business rules. A job consists of one or more steps, each representing a distinct phase of the workload. Simple jobs may have a single step, while complex processes can use multiple steps in sequence, such as importing data, validating it, and generating a final report. Within a chunk-oriented step, the ItemReader supplies data to the runtime one item at a time, from sources such as a database or file. The reader only retrieves the next item and does not need to know how it will be processed or persisted.The ItemProcessor receives each item and applies business rules, such as validation, transformation, classification, enrichment, or filtering. It may return a modified item or null if the item should be excluded from writing.The ItemWriter receives processed items and persists or exports them. Unlike the reader and processor, which handle items individually, the writer typically receives a group of items from the current chunk. This enables more efficient database or bulk operations. Jakarta Batch adds features around this pipeline to support enterprise workloads. The runtime manages chunk boundaries, transactions, checkpoints, execution status, failures, and restart behavior. Chunk size determines how much work is grouped before a write and checkpoint, making it a key parameter for balancing throughput, memory usage, database cost, and recovery. The core model is straightforward: Job → Step → Read → Process → Write → Repeat Jakarta Batch keeps the business pipeline simple while the runtime manages the execution mechanics needed for reliable, long-running data processing. The Sample: Customer Segmentation with Jakarta Batch This example demonstrates the Jakarta Batch model using an e-commerce customer segmentation scenario. Customers are assigned to tiers such as Bronze, Silver, Gold, and Platinum based on configurable spending thresholds. When thresholds change, the application reevaluates the customer base and updates only customers whose classification has changed. The full application includes MongoDB integration, a Jakarta Faces UI, a preview workflow, validation, and supporting services. The complete source code is available at https://github.com/soujava/mongodb-jakarta-batch. This section focuses on the classes directly involved in Jakarta Batch execution. Starting the Batch Job The application initiates the batch process through CustomerSegmentationService. Unlike the reader, processor, and writer, this class is not a batch artifact. Instead, it is an application service that retrieves Jakarta Batch’s JobOperator from BatchRuntime to start and monitor job executions. Java @ApplicationScoped public class CustomerSegmentationService { public static final String JOB_NAME = "customer-segmentation"; private volatile CustomerSegmentationPolicy currentPolicy; // initialization and status methods omitted public long start(CustomerSegmentationPolicy policy) { if (isRunning()) { throw new IllegalStateException( "A customer segmentation batch is already running"); } Properties parameters = new Properties(); parameters.setProperty( CustomerSegmentationPolicy.JOB_PARAMETER, policy.toJson()); long executionId = BatchRuntime.getJobOperator() .start(JOB_NAME, parameters); currentPolicy = policy; return executionId; } public boolean isRunning() { // implementation omitted } } The key API here is JobOperator, which Jakarta Batch provides as the interface for starting, stopping, restarting, and inspecting jobs. In this example, the segmentation policy is serialized into the job parameters to ensure each execution gets the correct business rules. Reading the Input The first batch artifact, CustomerItemReader, extends Jakarta Batch’s AbstractItemReader to implement a chunk-oriented reader. Java @Named("customerItemReader") @Dependent public class CustomerItemReader extends AbstractItemReader { @Inject private CustomerRepository customerRepository; private List<Customer> customers = List.of(); private int nextIndex; @Override public void open(Serializable checkpoint) { try (Stream<Customer> customerStream = customerRepository.findAll()) { customers = customerStream .sorted(Comparator.comparing(Customer::getId)) .toList(); } nextIndex = checkpoint instanceof Integer index ? index : 0; } @Override public Customer readItem() { if (nextIndex >= customers.size()) { return null; } return customers.get(nextIndex++); } @Override public Serializable checkpointInfo() { return nextIndex; } } These methods are part of the Jakarta Batch reader lifecycle defined by AbstractItemReader. open() prepares the reader and accepts a previous checkpoint if available. readItem() provides the next item to the runtime; returning null indicates there is no more input. checkpointInfo() reports the reader’s current position for checkpointing. For simplicity, this sample loads customers into memory. For larger workloads, the implementation might use pagination or a MongoDB cursor without changing the Jakarta Batch model. Processing Each Customer The next artifact implements Jakarta Batch’s ItemProcessor interface. Java @Named("customerTierProcessor") @Dependent public class CustomerTierProcessor implements ItemProcessor { @Inject @BatchProperty( name = CustomerSegmentationPolicy.JOB_PARAMETER) private String thresholdsJson; private CustomerSegmentationPolicy policy; @PostConstruct void initialize() { policy = CustomerSegmentationPolicy.fromJson( thresholdsJson); } @Override public Customer processItem(Object item) { if (!(item instanceof Customer customer)) { throw new IllegalArgumentException( "Expected a Customer item"); } CustomerTier calculatedTier = policy.tierFor(customer.getTotalSpent()); if (calculatedTier == customer.getTier()) { return null; } return Customer.builder() .id(customer.getId()) .name(customer.getName()) .totalSpent(customer.getTotalSpent()) .tier(calculatedTier) .build(); } } Here the Jakarta Batch contract is explicit: ItemProcessor defines processItem(). The runtime calls that method for every item produced by the reader. The processor applies the segmentation rule and either returns the transformed customer or null. Returning null has a specific meaning in Jakarta Batch: the item is filtered and does not continue to the writer. The @BatchProperty is also part of the Batch integration. It receives the thresholds property defined for this job execution, allowing the processor to reconstruct the CustomerSegmentationPolicy before processing begins. Writing the Results The final artifact extends AbstractItemWriter, Jakarta Batch’s base implementation for writing a chunk. Java @Named("customerItemWriter") @Dependent public class CustomerItemWriter extends AbstractItemWriter { @Inject private CustomerRepository customerRepository; @Override public void writeItems(List<Object> items) { List<Customer> customers = items.stream() .map(this::toCustomer) .toList(); customerRepository.saveAll(customers); } private Customer toCustomer(Object item) { if (item instanceof Customer customer) { return customer; } throw new IllegalArgumentException( "Expected a Customer item"); } } writeItems() is defined by the Jakarta Batch writer contract inherited from AbstractItemWriter. Unlike the processor, which receives one item at a time, the writer receives a collection of processed items. In this case, the collection contains only customers whose classification changed, as the processor has already filtered the others. At this point, the Java components of the pipeline are as follows: Plain Text CustomerItemReader extends AbstractItemReader ↓ CustomerTierProcessor implements ItemProcessor ↓ CustomerItemWriter extends AbstractItemWriter These types are what connect the application code to the Jakarta Batch runtime. Connecting the Artifacts With JSL The Java classes define the behavior, but Jakarta Batch requires explicit mapping of the reader, processor, and writer to each job. This orchestration is described in JSL: XML <?xml version="1.0" encoding="UTF-8"?> <job id="customer-segmentation" xmlns="https://jakarta.ee/xml/ns/jakartaee" version="2.0"> <step id="recalculate-customer-tiers"> <chunk item-count="20"> <reader ref="customerItemReader"/> <processor ref="customerTierProcessor"> <properties> <property name="thresholds" value="#{jobParameters['thresholds']}"/> </properties> </processor> <writer ref="customerItemWriter"/> </chunk> </step> </job> The ref values correspond directly to the names declared with @Named in the Java classes: Java @Named("customerItemReader") @Named("customerTierProcessor") @Named("customerItemWriter") The XML therefore tells the Jakarta Batch runtime: for this step, use this reader, then this processor, and finally this writer. It also maps the thresholds job parameter into the processor property. The item-count="20" sets the chunk size for this sample. Jakarta Batch coordinates reading and processing, periodically invoking the writer according to the chunk lifecycle and establishing transaction and checkpoint boundaries. The value 20 is for demonstration; real applications should tune chunk size based on processing cost, database behavior, transaction size, throughput, and recovery requirements. This structure is recommended for the article: present the class declaration first, then describe the lifecycle methods inherited from or required by Jakarta Batch. This approach helps the sample teach the API rather than simply presenting isolated methods. Conclusion Jakarta Batch is valuable because it transforms large-scale data processing into a structured execution model, eliminating the need for custom loops and ad hoc background logic. By separating reading, processing, and writing, and introducing runtime concepts such as jobs, steps, checkpoints, restartability, and chunk-oriented execution, it provides enterprise applications with a predictable approach to handling workloads involving thousands or millions of records. This allows implementations to focus on business logic, while the Batch runtime manages repetitive execution concerns, making the model easier to understand, optimize, and scale as workloads increase.

By Otavio Santana DZone Core CORE
Stop Paying a Model to Make Decisions You Already Made
Stop Paying a Model to Make Decisions You Already Made

If your team distributes AI development skills (as a Claude Code or Cursor plugin or a shared rules file), you own a catalog. That catalog covers things like: How to structure a serviceWhat must pass before a commitWhich internal library to use instead of rolling your own Skills load cheaply, thanks to progressive disclosure. Only the name and one-line description sit in context until the description matches the work. So when token spend climbs, the intuitive read is context bloat, and the obvious lever is tuning descriptions to fire less. That work is worthwhile as it improves routing quality. But it addresses only one of the two ways a skill costs money. The other is procedural amplification: a skill that encodes a fixed procedure as prose, so the model reconstructs it, deliberating, calling tools, checking results, on every invocation. That doesn’t show up as a large context payload. It shows up as extra inference round trips, and an aggregate cost dashboard cannot see it. Predictability does not prove a step should be automated. It identifies places where you should ask whether you’re paying a model to make a decision the system has effectively already made. A few terms come up more than once below, defined here so they don’t slow you down later: TermMeaningSkillA packaged set of instructions Claude Code loads only when the task matches it.OTelShort for OpenTelemetry, the open standard Claude Code uses to report what it did and how much it cost.API request / round tripOne call to the model and its response. This article counts these to measure cost.HookA small script Claude Code runs automatically at a specific moment, such as right after every tool call.p50 / p90 / p95 / p99Percentile timing. p50 is the typical (median) call; p99 is close to the worst case you’ll see.AsyncA setting that lets a hook run in the background instead of making the agent wait for it to finish.NDJSONOne JSON record per line in a file. Simple to append to and to stream. Use the Vendor’s Telemetry With Small Customization I’ll correct a claim I believed myself and have seen repeated: OpenTelemetry only provides aggregate counters, so you have to build your own event pipeline. That’s false, and starting from scratch costs you a week. Claude Code’s OTel export includes events, and they’re richer than I expected: EventCarriesclaude_code.api_requestskill.name, cost_usd, input_tokens, output_tokens, duration_ms,event.sequence, prompt.idclaude_code.tool_resulttool_use_id, tool_name, success, duration_ms skill.name on the API request is the important one: inference round trips are already attributed to the skill that was active. cost_usd is documented as an estimate, not billing data. The One Field to Add Getting command detail out of the native exporter requires OTEL_LOG_TOOL_DETAILS=1, which attaches full_command, bash_command, and file_path to tool events. That’s exactly the payload you cannot centralize off a fleet of developer laptops: branch names, commit messages, customer identifiers, the occasional pasted token. So a small local hook fills one gap: a privacy-preserving shape of the command, joined to native telemetry on tool_use_id, which the docs describe as matching the tool_use_id passed to hooks, allowing correlation between OTel events and hook-captured data. Plain Text native OTel ──┐ │ api_request: skill.name, cost, tokens, round trips ├── join on tool_use_id ──▶ analysis │ tool_result: tool_use_id, tool_name, success local hook ──┘ The hook doesn’t rebuild the trace. It emits the join key plus the one thing native telemetry can’t safely give you. Anything OTel already reports (success, duration_ms, skill.name, token counts) is deliberately not duplicated. Two sources of truth for one field will eventually disagree, and you’ll trust the wrong one. The command-to-shape transform itself is a per-binary allowlist, not a secret detector, and it has a sharp edge worth knowing: The transform: git commit -m "fix auth for acme" becomes git commit -m <ARG>.The edge case: a value-taking flag like --token abc123 leaks its argument as a bare token unless you explicitly track which flags consume the next word. Get the allowlist wrong, and the safety argument for the whole pipeline goes with it. What the Hook Costs (Measured) Tool hooks run on the critical path, so I benchmarked one: 1,000 synthetic PostToolUse payloads through the real script, mixed across Bash/Read/Edit/Write/Grep/Skill. Plain Text p50 p90 p95 p99 max mean subprocess 24.28ms 26.91ms 27.91ms 32.31ms 52.42ms 24.93ms in-process 0.08ms 0.15ms 0.16ms 0.20ms 0.35ms 0.10ms macOS, arm64, 10 cores, Python 3.9.6, n=1000 after 25 discarded warmups. The interesting part isn’t the total; it’s the split. About 24.2 ms of that 24.3 ms p50 is Python interpreter startup. The actual work — scrubbing, serializing, appending — costs 0.08 ms. That points the fix somewhere I would not have guessed: Don’t optimize the scrubber. It’s already three orders of magnitude below the launch cost.Do fix how the hook is launched. Register it async: true so nobody waits, or make it a thin client to a long-lived local collector so you pay startup once per session instead of per tool call. Had I estimated instead of measured, I’d have guessed low single-digit milliseconds and gone tuning the scrubbing code. Wrong target entirely, which is the argument for measuring rather than estimating, in miniature. Event records came in at 364 bytes each. One engineer at roughly 40 sessions a week and 60 tool calls a session generates about 2,400 records, under 1 MB of raw NDJSON weekly. It stays small-data at any plausible team size, because the enrichment record only carries what native OTel doesn’t. Measuring Your Own Catalog Skills have two token costs that behave differently, and knowing which one dominates for you tells you where to look first: Cost componentPaid whenScales withDescriptionEvery request, every session, whether the skill fires or notCatalog sizeBody (procedural amplification)On invocation, then re-sent on each subsequent request in the sessionHow often the skill fires, and how long sessions run Measuring the 25 skills in Claude Code’s official plugin marketplace, a catalog anyone can install and re-measure, gives: Plain Text body chars: min 989 median 11,395 max 32,625 (33x spread) desc chars: min 105 median 357 max 906 always-resident total (all 25 descriptions): 9,703 chars Read those two numbers against each other. Carrying every description in the catalog costs less than a third of one large body. The single biggest skill is roughly 3.4 times the entire always-resident footprint, and it’s charged again on every request for the rest of any session that invokes it. That’s the shape that makes procedural amplification worth hunting. If your always-resident total instead dwarfs your median body, the lever really is catalog size and description tuning. You can stop there. (These are characters, not tokens. The ratio shifts with how much code versus prose a body contains. Count with the real tokenizer, messages.count_tokens, which is free and rate-limited only by requests per minute; don’t reuse a chars/4 estimate or another vendor’s tokenizer, and don’t reuse an old Claude count either: the tokenizer used by Claude 4.7 and later yields roughly 30% more tokens for the same text than earlier models.) Why This Is Worth Measuring at All The subject isn’t token optimization. It’s finding misplaced probabilistic computation: places where the system already knows the next operation and is paying a model to rediscover it. If the next operation is…It belongs in…Genuinely uncertainA skill, policy and judgment, model reasoningAlready determinedA script, deterministic execution A skill that spells out a fixed procedure in prose has put that boundary in the wrong place. The industry is converging on the same line from the runtime side: programmatic tool calling exists precisely to keep deterministic multi-tool sequences out of the inference loop. Key Takeaways Aggregate cost accounting answers the question you already knew to ask. Event-level behavioral traces let you find the one you didn’t, and in Claude Code, most of that stream already exists, attributed to the skill, waiting on one privacy-preserving field to become useful. Procedural amplification leaves no trace in a token counter, raises no error, and turns no dashboard red. It’s visible in the order of the calls, and in how many round trips it takes to get through them. Measured, not estimated: Hook overhead is about 24 ms per call, almost entirely interpreter startup, not the redaction logic itself.Measured, not estimated: In one real catalog, a single skill body costs 3.4× what the entire always-resident description set costs.Next step: Measure by the actual stretch of work, not by an arbitrary number of calls. Group everything that happens while one skill is active, however long that turns out to be, rather than picking a fixed count upfront. When one model turn fires off several tool calls at once, count that as the single decision it was, not several. And don't trim the long, expensive stretches out of the data before analyzing it. Those are usually the exact cases worth finding.

By Amith Reddy Ravuru
Beyond Batch: Engineering Enterprise Systems for Real-Time Decisioning
Beyond Batch: Engineering Enterprise Systems for Real-Time Decisioning

The Batch Processing Problem Batch processing isn't inherently a disadvantage. It becomes a problem when the business needs a decision now, but the architecture was designed to make that information available later. Picture an enterprise system processing millions of customer interactions. Transactions land across multiple systems throughout the day. Every few hours, a scheduled job extracts the data, transforms it, updates another system, and eventually makes it available downstream. This works fine — until the business asks: "Why can't we react to this the moment it happens?" A suspicious transaction. A changed preference. A completed payment. A signed document. A submitted service request. An account state change. In large enterprise environments, I've watched a fairly consistent pattern play out. Teams initially focus on throughput and infrastructure capacity — can the pipeline handle the volume, can it finish the batch window in time? As the systems mature, the harder questions shift elsewhere entirely: who owns a given event, how failures get recovered, how schemas evolve without breaking consumers nobody remembers exist, and — most importantly — what the actual business impact is when a consumer falls behind. Running the batch more frequently doesn't answer any of those questions. Eventually, the architecture itself has to change. Event-Driven Doesn't Mean "Install Kafka" This is where most transformations quietly stall. A common pattern looks like: batch system → add Kafka → the same tightly coupled design underneath. The organization now calls itself event-driven, but nothing structural has actually changed. Real event-driven architecture requires rethinking state ownership, service boundaries, data contracts, failure handling, consistency assumptions, observability, and operational responsibility — not just swapping the transport layer. Plain Text Business Action | v Producer Service | v Event Backbone / | \ v v v Risk Customer Analytics Svc Svc Svc The producer shouldn't need to know or care who's downstream. That's the real architectural benefit — a producer that stays ignorant of its consumers is what genuine decoupling looks like. If the producer still has to know which five systems need to be updated and in what order, you haven't built an event-driven system — you've built a batch job that happens to run on Kafka. Business Events Are Not Technical Messages There's an important distinction between commands and events. A command — UpdateCustomerProfile, SendNotification — says do something. An event — PaymentAuthorized, DocumentSigned — says something happened. Well-designed events represent durable business facts, not implementation instructions. Publish PaymentAuthorized, and Fraud Detection, Notifications, Analytics, Accounting, and Audit can all react independently, without the producer orchestrating any of them. That's the difference between an event-driven system and a batch system wearing a streaming costume. The Hard Problems Start After the First Event Duplicates Most messaging systems guarantee at-least-once delivery, so PaymentAuthorized may legitimately arrive twice. The customer shouldn't be charged twice. Idempotency — via event IDs, business transaction IDs, or a processed-event store — isn't optional polish. Duplicate delivery is a normal condition in a distributed system, not an edge case you occasionally trip over. Ordering If AccountClosed is processed before AccountCreated ever arrives, the consumer ends up holding a state that shouldn't be able to exist — an account that's closed but was never opened. The instinct is to enforce global ordering everywhere, but that kills scalability. The better question is narrower: what actually needs to be ordered? Usually it's events for the same business entity — the same customer, the same account — not the entire enterprise-wide stream. Schema Evolution An event schema gains a new field six months after a consumer was deployed against the old one. Does it break? Backward compatibility, schema registries, and contract testing matter here because events tend to outlive the applications that created them. Treat event contracts like governed APIs, not like internal implementation details nobody needs to track. Failure Don't retry forever. A sane strategy escalates in stages: an initial attempt, then a short retry, then backoff, then a dedicated retry queue, then a dead-letter queue for anything that still hasn't succeeded, then manual investigation or replay. A malformed or logically invalid "poison" event shouldn't be allowed to block the pipeline indefinitely just because it keeps failing the same way. Worth watching closely: retry count, dead-letter volume, consumer failure rate, and the age of the oldest unprocessed event. Eventual Consistency Changes How Teams Think In a synchronous system, an update and its visibility happen together — you write, you read back the new value, done. In an asynchronous architecture, that guarantee disappears. One consumer might reflect a change in twenty milliseconds; another might take two seconds; a third might be temporarily unavailable and catch up later. Different systems can legitimately hold different states for a period of time, and that isn't automatically a defect. The real architectural question is: how stale can this information safely become? Fraud decisioning tolerates almost none — a few hundred milliseconds of staleness can be the difference between catching and missing something. Marketing analytics can tolerate a great deal more. Audit cares more about completeness than about speed. This needs to be decided per business function, not applied as one blanket policy across the platform. It's also worth being honest about what "real-time" actually means in practice. A system that processes an event in milliseconds isn't meaningfully real-time if the downstream systems that act on that event take minutes to reflect the result. I've seen teams celebrate a fast event pipeline while the actual customer-facing decision — the offer shown, the risk flag raised — still lagged well behind because a downstream dependency hadn't caught up. Real-time decisioning has to be measured end-to-end, at the point where the business decision is made, not just at the point where the event was published. Migrating Off Legacy Without a Big-Bang Cutover Ripping out a legacy system in one motion rarely goes well. A more workable path is incremental: capture changes from the legacy system as events, route them through the event backbone, and let new and existing systems consume from the same stream during the transition. Plain Text Legacy System | v Change / Event Capture | v Event Backbone / | \ v v v New New Existing Svc Svc Systems The transactional outbox pattern is useful here: write the event to an outbox table in the same database transaction as the business update, then publish from the outbox separately. That avoids the classic dual-write problem, where the database commit succeeds but the event publish fails, silently leaving downstream systems out of sync. Change Data Capture can also help expose changes from a legacy system as a migration bridge. But it's worth being deliberate about this: a raw database row change is not automatically a well-designed business event. CDC tells you a row changed; it doesn't tell you why, or whether that change represents something a downstream consumer should actually care about. Treating every CDC record as a business event is one of the more common ways these migrations end up producing noisy, low-value streams. Observability Has to Be Designed In, Not Added Later A customer says: "My transaction disappeared." Where do you look, across a chain of services and events? You need correlation IDs, trace IDs, event IDs, business transaction IDs, timestamps, producer identity, and schema versions threaded through everything — plus the standard infrastructure metrics: consumer lag, event age, processing latency, retry rates, dead-letter volume, error rate. But infrastructure telemetry on its own isn't enough. Knowing "consumer lag is 12,000" is far less useful than knowing "12,000 customer transactions are currently delayed." That translation — from technical signal to business impact — is what tends to separate a platform that's merely instrumented from one that's genuinely observable. It's also usually the gap that shows up first when something goes wrong in production: the engineering team sees a metric, and it takes real effort to connect that metric to what a customer or a business stakeholder is actually experiencing. Security and Governance Have to Follow the Data Event-driven architecture multiplies how much data moves around a system, so security has to travel with the data rather than sit only at the application perimeter. That means clear authentication and authorization for who can publish and consume which topics, encryption both in transit and at rest, discipline about not routinely copying sensitive or personal information into every event just because it's convenient, defined retention policies, and clear auditability of who produced what and when. When Not to Use Event-Driven Architecture Don't adopt EDA because it's fashionable. A synchronous API is often the better choice when immediate request-response is required, the workflow is simple, only one system needs the result, or strong immediate consistency is essential. Batch remains entirely appropriate for monthly statements, historical reporting, bulk reconciliation, archival, and much regulatory reporting. The mature position isn't "everything must become event-driven." It's choosing synchronous, asynchronous, and batch patterns based on what the business actually requires — and having a clear answer for why. A Practical Decision Framework Before converting a workload, it's worth asking a short set of questions: Does the business genuinely require lower latency — or would nobody notice the difference between seconds and hours?Do multiple independent consumers need the same business change? If so, event-driven design becomes attractive.Can the business tolerate eventual consistency? If not, the workflow needs closer examination before proceeding.Can the organization actually operate distributed, asynchronous systems — with the observability, on-call practices, and schema governance that requires?What happens when one component fails? If the design can't answer that clearly before production, it isn't ready for production. A Reference Architecture Plain Text ┌──────────────┐ │ Channels │ └───────┬──────┘ │ v ┌──────────────┐ │ API / Domain │ │ Services │ └───────┬──────┘ │ Business Events │ v ┌────────────────────────┐ │ Event Backbone │ └────────────────────────┘ │ │ │ ┌────┘ │ └────┐ v v v Decisioning Notifications Analytics │ │ │ v v v Data Store Data Store Data Store ──── Observability ──── ────── Security ─────── ───── Governance ────── Observability, security, and governance aren't downstream services bolted onto the diagram — they span the whole architecture, or they don't really work. Conclusion The real transformation isn't batch → Kafka. It's delayed processing → continuous business awareness, and central orchestration → autonomous consumers responding to business facts. That shift comes at a cost: more distribution, more asynchronous behavior, more operational complexity, more governance overhead. So the goal was never to produce more events. The goal is systems capable of making timely, reliable decisions — while staying understandable and operable when, inevitably, something fails.

By Prem Kumar Gadhanki
The Hidden Production Risks of Third-Party SDKs
The Hidden Production Risks of Third-Party SDKs

Most modern applications do not function completely independently. For example, analytics, payment processing, user authentication, customer support, testing new features (experimentation), monitoring the app's performance, advertising, etc., are typically provided as third-party SDKs that enable those functions in your app. Using an SDK has its benefits; you don't have to build an entire piece of functionality yourself. When using an SDK, developers can download the software library, call the initialization method, then begin calling the API methods of the SDK to use its functionality. JavaScript import { analytics } from "third-party-sdk"; analytics.track("checkout_started", { productId: "123" }); In some cases, the amount of code required to add this type of functionality can be as little as a handful of lines of code. If the same functionality was built "from scratch", the time required could potentially be several weeks. However, with great convenience comes hidden complexity. The moment you allow a third-party SDK to run in your app, how well it performs and works (performance, reliability, security, and user experience) depends on what amounts to "someone else" doing something to your app. That is why third-party SDKs are important dependencies that affect how well your app will perform during production hours, instead of just being another library or module to include. SDKs Can Quietly Affect Performance The biggest reason front-end SDKs will show performance issues is that they typically run directly in your web browser. If you install a typical analytics SDK, it adds to your front-end bundle, it loads on page init, and then registers event handlers, makes requests over the internet, etc., as soon as there are interactions with your app. Although one SDK alone has little effect, if you use multiple SDKs for analytics, experimentation, customer service, session replay, ad tracking, and monitoring, your users will likely notice a difference. Therefore, teams need to evaluate whether individual "acceptable" SDK costs can compound into overall user-perceived degradation. In addition to measuring the cost of each SDK individually, teams need to look at the overall cost of loading the SDK(s), which can include: Bundled sizeTime to initializeNetwork requestsActivity on main threadOverall impact on Core Web Vitals This cost can be reduced by loading non-critical SDKs asynchronously or after the main application experience has loaded. A Third-Party Failure Can Become Your Failure Consider an application that will render the main page after initializing a recommendation SDK. JavaScript await recommendationSDK.initialize(); renderApplication(); If the third-party service has an issue, your application's overall appearance may slow or become unavailable, even if your backend is functioning properly. Thus creating unneeded coupling. Generally speaking, non-essential third-party services should be allowed to fail without affecting the primary user experience. For example, if you're unable to receive recommended products, you should still be able to browse through products; if analytics are failing, checkout should still function as normal; and if a support widget is unable to load, all other aspects of the webpage should continue to function normally. Applications should define clear fallback behavior for every external dependency. Timeouts are also important. Waiting indefinitely for a third-party service can turn a small external outage into a much larger product incident. SDK Updates Can Change Production Behavior Engineers typically spend considerable time evaluating large-scale framework updates; however, they may be less concerned about small third-party dependencies that make up much of their application codebase. This could potentially lead to issues. An SDK update can change how an application initializes, the format for making requests, which browsers an application supports, the default configuration, how data is stored, or how much JavaScript is downloaded during each session. Even if the public API hasn't changed, runtime behavior may still differ based on previous SDK versions. Therefore, dependency upgrades should follow standard engineering controls such as version pinning where applicable; automated testing; dependency review; and gradual deployment. The idea of automatically allowing all new SDK releases into production just because they have been classified as minor will create additional risk. Third-Party Code Expands the Security Boundary Every new SDK you add to your app will be a larger portion of all code making up the system. Browser apps make this especially important when SDKs can access page content, browser storage, cookies, user interaction, or even application data. You should know exactly which pieces of information will go out to an outside party. As an example, sending off an entire object to an analytics SDK could provide more information than was ever intended: JavaScript analytics.track("profile_updated", user); Some of the fields in the 'user' object might never have been intended for analytics. A safer approach is to explicitly select the information required for the event. JavaScript analytics.track("profile_updated", { accountType: user.accountType }); You'd be better off sending only the data you need for each specific event. The principle is simple: third-party integrations should receive only the data they actually need. SDKs Can Create Hidden Runtime Conflicts Not all third-party SDKs run independently. In addition to other actions such as modifying a browser's global objects, registering event handlers, intercepting web requests, and manipulating the DOM, third-party SDKs may also create new dependency conflicts that are incompatible with your current application code. Because of their nature, these issues can be difficult to reproduce because they typically depend on specific conditions (such as browser type, user environment, feature flags/feature toggle configuration, etc.) that cause them to occur only under very specific circumstances. Another reason why you should track correlation of failures to your integrations during production time is due to this. Also, when possible, initialize third-party SDKs in an isolated manner so that an error in initializing one service does not bring down the rest of the application. JavaScript try { await supportSDK.initialize(); } catch (error) { logError("Support SDK initialization failed", error); } If your optional service fails to start, then your application continues. Have an Exit Strategy The other, quite surprising, issue you might have when using an SDK is the difficulty in removing it. When there are numerous API calls in multiple layers of your app, making changes to which vendor you use as a service provider becomes extremely expensive. In this case, teams may want to develop their own internal abstraction layer on top of the external SDK. JavaScript tracking.track("checkout_started", data); Your application interacts with an internal 'tracking' interface, and then the internal tracking layer will interact with the external SDK. You still have a dependency on the vendor, but now all vendor-specific APIs are abstracted out of your codebase. Testing also becomes simpler, and you can easily add validation, filtering, error handling, and fallback logic. Monitor SDKs Like Production Dependencies Integrations with third-party tools should look similar in your observability dashboard as your internal services. Understanding when/why an SDK will fail; how long initialization takes; whether requests are timing out; and which specific integration(s) cause frontend errors/performance regressions helps teams understand when they have a problem. It is also beneficial to understand what feature of your application depends on each provider. This type of information greatly assists during an incident by providing a clear yes/no answer to an important question: Can I disable this integration and still run my core product? For critical integrations, the answer should already exist before an outage occurs. Conclusion Third-party SDKs are useful because they enable engineering teams to get things done in less time than would be required if the team had to build capability again, which has been developed by others who specialize in that area of development. However, when you add a new SDK to your project, you've added a new production dependency. This dependency can negatively affect your application's performance, reveal information about your application, break at unpredictable times, change with each upgrade, and ultimately become very hard to remove as it spreads across your codebase. Our objective is not to eliminate third-party SDKs. Our goal is to intentionally incorporate third-party SDKs into the project. Track how much performance is affected, track what amount of data is transmitted back to the provider, prevent failures from spreading through isolation, maintain control over upgrades, track the use of the service, and do everything possible to prevent tightly coupling core functionality to a service that the application does not control. A third-party SDK may take only a few minutes to install, but its production impact can last for years.

By Satyam Nikhra
Stop Blaming Executor Memory: The Real Reasons Your Spark Jobs Are Slow
Stop Blaming Executor Memory: The Real Reasons Your Spark Jobs Are Slow

After a decade of building and debugging large-scale data pipelines across financial services, payments processing, and analytics platforms, I can tell you that almost every slow Spark job I've investigated had the same root cause — and it wasn't the one the team thought it was. The default response when a Spark job is slow is to add more executor memory, increase the number of executors, or bump spark.sql.shuffle.partitions. Sometimes that helps. Usually it doesn't. What I've found, consistently, is that the real problems are structural — a join strategy mismatch that silently multiplies your intermediate dataset by ten times, a single slow task on a degraded node that holds an entire stage hostage, or a decrypt chain that re-reads source data six times when it only needed to read it once. This article is organized around five patterns I keep seeing across teams. Each one looks different on the surface but traces back to a misunderstanding of how Spark actually executes your code. For each pattern, I'll describe what it looks like, when it bites you, the failure mode, and how to fix it. Pattern 1: The OR Join That Quietly Multiplies Your Data What It Looks Like A join condition with an OR clause. Usually introduced when a business requirement adds a secondary matching rule — match on primary card number, or if the transaction is a virtual card transaction, match on the underlying physical PAN. The SQL looks reasonable. The engineer tests it on a sample, and it returns the right rows. When It Bites You At scale. With 100 million transaction rows and 50 million account rows, this query starts running for hours. The output size is also wrong — much larger than expected before DISTINCT trims it down. The Failure Mode Spark cannot use a hash join or sort-merge join when the join condition contains OR. It falls back to BroadcastNestedLoopJoin — for every row in the left table, scan every row in the right table. That's O(n x m). On real datasets, this produces an intermediate result in the hundreds of GB before any downstream filter runs. I've watched a pipeline that should produce 8 GB of output generate 400 GB of intermediate data because of exactly this pattern, taking a 20-minute job to 4 hours. You can verify this in 30 seconds: run df.explain(formatted) and look for BroadcastNestedLoopJoin in the physical plan. If you see it on a join involving any table over a few million rows, it's almost certainly unintentional. The Fix Split the join into two equi-join legs and UNION ALL the results: SQL -- Leg 1: primary match (equi-join — uses SortMergeJoin or BroadcastHashJoin) SELECT txn.*, acct.* FROM transactions txn JOIN accounts acct ON txn.card_number = acct.card_number UNION ALL -- Leg 2: fallback match, filtered scope only SELECT txn.*, acct.* FROM transactions txn JOIN accounts acct ON txn.fpan = acct.physical_pan WHERE txn.transaction_type = 'VIRTUAL' Each leg is a proper equi-join. Apply DISTINCT at the end to deduplicate rows that matched both. The performance difference is routinely an order of magnitude. Pattern 2: The Straggler Task That Nobody Notices Until It's Too Late What It Looks Like A stage that should take 10 minutes takes 3 hours. The Spark UI shows nearly all tasks completed quickly. One or two tasks are still running with a disproportionately long duration. When It Bites You Jobs running on shared YARN or cloud infrastructure where any node can have a bad disk, a noisy neighbor, or degraded network throughput. Also common in stages that call external services per partition — one slow API response can cause a single partition's tasks to take 100x longer than the others. The Failure Mode A stage doesn't complete until the last task completes. Not the median. Not p95. The absolute last one. If 2,200 tasks finish in under 2 minutes and one takes 3 hours and 7 minutes, the stage takes 3 hours and 7 minutes. The other 2,199 executors sit idle. This is the straggler problem, and it's distinct from data skew. The diagnostic: in the Stage detail view, check the task duration distribution. If MAX is dramatically higher than p99, that's a straggler (hardware or external service issue). If p75 is already much higher than p50, that's skew (data distribution issue). They require different fixes, and many teams treat them identically. The Fix For stragglers caused by degraded infrastructure, enable Spark speculation: Properties files spark.speculation=true spark.speculation.multiplier=3 # task must be 3x slower than median spark.speculation.quantile=0.9 # wait for 90% completion before speculating Speculation re-launches slow tasks on a different executor and uses whichever copy finishes first. The caveat: don't use this on stages that write to non-idempotent sinks. For read-heavy or compute-heavy stages — including external decryption calls — it's often the single most impactful config change you can make. Pattern 3: The df.rdd Decrypt Chain That Recomputes Everything Six Times What It Looks Like A pipeline that calls an external encryption or decryption service per record, implemented as a series of df.rdd.mapPartitions() calls, one per column that needs to be processed. When It Bites You When you have multiple columns to decrypt. Each .rdd call creates a new computation starting from the original DataFrame — Spark re-reads from source, re-executes all upstream joins and filters, and then runs the decryption for that column. With six columns to decrypt, you're doing that six times. The Failure Mode Two distinct sub-problems compound each other. First, going to RDD bypasses Catalyst entirely — no predicate pushdown, no column pruning, no Tungsten execution. Second, without a persist checkpoint before the chain, every decrypt call lineages all the way back to the source. I've seen this double the runtime of a job compared to the same pipeline with a single persist() before the decrypt chain. On top of that, the external call latency per partition is dominated by the number of HTTP round trips, not the payload size. Cutting your batch size in half doubles your request count and roughly doubles your wall-clock time for that stage. Most teams set an initial batch size and never revisit it. The Fix Two changes, applied together: Persist the input DataFrame before starting the decrypt chain. This means the join and filter logic runs once, and each decrypt call reads from the cached result.Increase the batch size for external calls. Test at several sizes — going from 20,000 to 40,000 records per batch often cuts stage time by 30-50% with no change to correctness. Scala val base = rawDf.filter(...).join(key1, ...).persist(StorageLevel.MEMORY_AND_DISK) val step1 = decryptColumn(base, secret1) // reads from cache val step2 = decryptColumn(step1, secret2) // reads from cache val step3 = decryptColumn(step2, secret3) // reads from cache Without persist, step2 re-executes everything step1 did from source. With persist, each step reads from the in-memory result of the previous. Pattern 4: The shuffle.partitions Setting That Nobody Updates What It Looks Like A job that works fine in staging — where data volumes are 10% of production — but runs slowly, spills to disk, or produces thousands of tiny output files in production. When It Bites You When the default spark.sql.shuffle.partitions=200 is left unchanged. 200 partitions made sense as a default for medium datasets but is almost always wrong at production scale — either too few (huge partitions, memory pressure) or too many (tiny partitions, scheduling overhead, small files problem). The Failure Mode Too few partitions means each executor handles a disproportionately large chunk of data. With 200 partitions on a 1 TB shuffle, each partition is 5 GB. That will spill to disk. Too many partitions means thousands of 1 MB tasks — the scheduling overhead becomes significant, and your output has thousands of tiny files that hurt downstream readers. With Adaptive Query Execution (AQE) enabled in Spark 3.2+, this problem largely manages itself. AQE merges small post-shuffle partitions automatically and can handle modest skew. But AQE can't help if it's disabled, and it can't fix the upstream causes of extreme skew. The Fix Enable AQE if you're on Spark 3.2+: Properties files spark.sql.adaptive.enabled=true spark.sql.adaptive.coalescePartitions.enabled=true spark.sql.adaptive.skewJoin.enabled=true If you need to set shuffle.partitions manually, target roughly 128-256 MB per partition post-shuffle. For a 500 GB shuffle, that means 2,000-4,000 partitions. Set it high and let AQE coalesce down — that's cheaper than setting it low and getting OOM errors. Pattern 5: The Incremental Job That Degrades Silently Over Time What It Looks Like A job that runs in 15 minutes when first deployed and runs in 4 hours six months later. No code changes. No obvious data quality issues. The team attributes it to data growth. When It Bites You When the job fails a few times in a row, and the recovery accumulates multiple windows' worth of data. Or when the watermark logic was designed for small windows but nobody anticipated that the underlying join tables would grow significantly. The Failure Mode Two separate causes, often confused. First, if the watermark is a single timestamp and the job has been failing, recovery runs can accumulate large backlogs. A job that normally processes 2 hours of data may need to process 48 hours on first successful recovery, with no change to the resource configuration. Second, growth in reference data (like an accounts table or lookup table used in a join) increases the size of every run regardless of whether the incremental input grew. I've seen a 30-minute job become a 3-hour job purely because the accounts table grew from 10 million rows to 80 million rows over 18 months, while the OR join condition (see Pattern 1) meant that growth was amplified into the intermediate result. The Fix Two design principles that pay off over the lifetime of the pipeline: Track processed partitions explicitly rather than using a single timestamp watermark. This makes recovery granular — you can replay specific missing partitions without re-processing everything after them.Add a fast-path no-op check before initializing the full Spark session. Check whether any new partitions exist first. A 5-second check that exits early is much better than a 2-minute executor startup that discovers there's nothing to process. For the reference table growth problem: if your lookup table grows significantly, revisit whether it can be broadcast (small enough to fit in executor memory) or whether the join itself needs to be redesigned. Quick Diagnostic Reference Use this table to map what you observe in the Spark UI to the likely pattern and first action to take: WHat you observeLikely patternconfirm withfirst action MAX task duration >> p99 Straggler (Pattern 2) Task timeline in Stage UI Enable spark.speculation p75 >> p50 task duration Data skew Input bytes per task Repartition on join key; AQE skewJoin BroadcastNestedLoopJoin in explain() OR join (Pattern 1) df.explain( formatted) Rewrite as UNION of equi-joins Stage runtime grows week on week; no code change Incremental accumulation or reference table growth (Pattern 5) Input bytes trend in History Server Audit watermark logic; check reference table size OOM errors or heavy disk spill Too few shuffle partitions (Pattern 4) Spill metrics in Stage UI Enable AQE or increase shuffle.partitions The Common Thread Every pattern here traces back to the same underlying issue: Spark is executing something different from what the engineer intended. The OR join was intended as a flexible matching rule; Spark turned it into a nested loop. The decrypt chain was intended as six independent transformations; Spark turned it into six full re-reads of source data. The incremental job was intended to process one window of data; without proper watermark design, it occasionally processes twelve. The Spark UI has everything you need to see this — task distribution, input and output sizes, physical plans, spill metrics. Most teams open it when something breaks and close it once they find the obvious error. Opening it proactively, forming a hypothesis, and then confirming or refuting it in the metrics is the practice that separates engineers who consistently improve pipeline performance from those who add executor memory and hope for the best. The mistake isn't choosing the wrong config. It's not understanding what Spark is actually doing with your code.

By Swaminathan Sethuraman
Understand the Sidecar Pattern by Deploying n8n to AWS Fargate
Understand the Sidecar Pattern by Deploying n8n to AWS Fargate

A sidecar is a container that runs alongside another container as part of the same deployment unit. Just because two containers are in the same cluster or deployed around the same time doesn't make one a sidecar. There are two things that make a sidecar. First is that they share a network namespace, so they can reach each other over localhost rather than a network address. Second, they share a lifecycle. This means that they are created together, scaled together, and by default torn down together. Neither container has an existence independent of the other. The problem it solves is giving a specific concern its own boundary. For example, it can have its own filesystem, its own memory space, and often its own permissions or dependency set, without giving up the simplicity of deploying and operating one unit. You get isolation without paying for the operational overhead of running and coordinating a fully separate service. The test that defines the pattern across all of these is this: does it live and die with its partner container as one unit of deployment? If yes, it's a sidecar. If you have to reach it by hostname, through service discovery, or via a queue, it isn't one anymore. That is a separate service that happens to sit next to the first. That test matters because two adjacent patterns get called "sidecar" when they aren't: Decoupled worker/microservice. A separately deployed container, reached over the network, scaled on its own. A web application offloading work to Celery workers via Redis is a common instance of this: the app enqueues a job (send this signup email), a pool of workers pulls jobs off the queue independently, and neither side shares a network namespace or a lifecycle with the other. The workers scale on queue depth, not on how many web replicas are running, and a web app restart doesn't take queued or in-flight jobs down with it. n8n has its own version of the same shape: "queue mode," where a main node accepts webhooks and separate worker nodes pull jobs off a Redis queue. It's tempting to call either of these a sidecar relationship since the worker and the web app do feel paired, but neither qualifies: they don't share a deployment unit, and killing one doesn't touch the other.Ambassador/adapter. A container that proxies or translates traffic on its parent's behalf, like the Envoy example above, is actually this, more precisely. Structurally it's still a sidecar; it just gets a more specific name for what it does. Using n8n to Understand It What n8n Is n8n is a workflow automation platform like Zapier, but self-hostable and node-based rather than form-based. A handful of components make up a running instance: The editor/UI, where workflows are built visually as a graph of nodes.The main process, which serves that UI, listens for webhooks, and orchestrates workflow execution. The workflow execution decides what runs next, passing data between nodes and recording results.Nodes, the individual units of a workflow: trigger nodes (a webhook arrives, a schedule fires), action nodes (call an API, write to a database, send an email), and the Code node. The code node lets you drop in arbitrary JavaScript or Python to transform data however the built-in nodes can't. The code node is relevant in this article. The database, where workflow definitions, credentials, and execution history persist. In this article, Postgres is used. For most of what n8n does, the main process is the only thing doing work: routing a webhook, calling an API, writing a database row. The exception is the Code node, and that exception is the whole reason task runners exist. The Task Runner Feature and Its Use Case By default, a Code node's JavaScript or Python executes inside n8n's main process. This main process holds the database connection, the encryption key, and every credential stored in every workflow you've built. That's fine for trusted, well-understood scripts. It becomes a real problem the moment the code in that node is untrusted, third-party, or arbitrary enough that you can't fully audit it before it runs. By the way, that is how most Code nodes are used in practice. Task runners exist to solve exactly that use case: run Code node logic somewhere the main process's credentials and connections aren't reachable from it, without turning "write some JavaScript to reshape this JSON" into a separately deployed microservice every time. Going Deep on the Task Runner Feature n8n ships two modes for this: Internal mode (the default) runs Code nodes inline, in-process. No isolation. This is the fastest to set up, but the weakest boundary.External mode moves execution into a separate runner process entirely. That process connects back to the main n8n instance over a broker (an authenticated connection the main process listens on) and receives individual tasks to execute rather than having any standing access to n8n's internals. The runner never touches the database connection, the encryption key, or stored credentials directly; it only ever sees the specific input data for the task it's been handed. External mode goes further than just "a different process," too. The runner's own configuration (the n8n-task-runners.json file built in Phase 4) sets explicit allowlists — which environment variables the runner process can see at all, and which JavaScript built-ins or Python modules it's permitted to import, standard library and third-party tracked separately. So the boundary isn't just "different memory space," it's "different memory space, plus a declared, auditable list of exactly what this process is allowed to touch." That's a specific concern (arbitrary code execution) given its own boundary, without turning it into a fully independent service you have to deploy, discover, and monitor separately. It's the sidecar problem, stated exactly: external mode gives you the isolation; running the external runner as its own container in the same task definition is what makes that isolation a sidecar rather than just a separate process sharing a machine. Why This Needs to Scale Independently and Why "In the Same Container" Isn't Enough Most n8n deployment guides run n8n with task runners in internal mode, or with the external runner living inside the same container as the main process. For example, you will see guides about deploying n8n on a single EC2 instance, Render, DigitalOcean, or any platform's basic tier. That gets you the process isolation, which solves the security half of the problem. It doesn't solve the other half, which is that a runner sharing a container with the app can't be scaled, resourced, or restarted independently of it. That stops mattering the moment Code-node execution becomes the actual bottleneck rather than webhook handling or UI traffic. Imagine workflows doing heavy data transformation in Python, running numpy/pandas operations across large payloads, or executing many Code nodes concurrently. If the runner is bundled into the main container, giving it more CPU means giving the entire n8n instance more CPU, whether the UI and webhook layer need it or not. There's no way to say "the runner needs 2 more vCPUs, n8n itself is fine". Why AWS Fargate's Task Definition Is the Right Fit A Fargate task definition lets each container in the task carry its own CPU and memory reservation, its own health check, and its own essential flag governing what happens if it fails while still keeping every container in the task on one shared network interface. That's the sidecar promise made literal: isolation and independent resourcing for the runner, without losing the operational simplicity of one task, one deploy, one thing to scale as a unit when you do want to scale both together. The rest of this guide deploys exactly that: one Fargate task, two containers, wired together the way the definition above requires. Each infrastructure decision below gets tied back to a specific part of what's laid out here, so that by the end, the concept isn't something read once at the top, but it's something built. Prerequisites AWS account with billing enabledA domain you control, with DNS accessDocker installed locally, with docker buildx availableAWS CLI configured (aws configure) with permissions for ECR, ECS, RDS, ACM, and IAMThe runner image source (Dockerfile + n8n-task-runners.json) — built in Phase 4 Architecture Markdown User's Browser (HTTPS) | [Application Load Balancer] <- Certificate Manager (SSL Cert) | (Port 5678, HTTP internal) [ECS Fargate Task] |-- Container: n8n (main) <-- shared network namespace --> Container: n8n-runner (sidecar) | (Port 5432, PostgreSQL) [RDS PostgreSQL Database] The load balancer and RDS layers are ordinary AWS plumbing. The box in the middle is where the sidecar relationship actually lives. There is one task and two containers, each with its own resourcing. Phase 1: RDS PostgreSQL RDS Console → Create database → Standard create → Engine: PostgreSQLDB instance identifier: n8n-db. Master username: postgres. Generate and save a strong master password.Instance size: db.t4g.microStorage: 20 GB gp3, autoscaling on, max 100 GBConnectivity: the VPC you'll use throughout. Public access: No. New security group: n8n-db-sg, left empty for now.Additional configuration → Initial database name: n8n. Skip this and n8n fails on first connect with "database does not exist" — the DB instance identifier names the server, this field names the database inside it.Create, wait for "Available," copy the endpoint from Connectivity & security. Phase 2: ACM Certificate n8n requires HTTPS for webhooks to function Certificate Manager, in the same region you'll deploy the Load Balancer in → Request a public certificateDomain name: n8n.yourdomain.comValidation method: DNS validationCreate the CNAME record ACM provides at your registrar. If your registrar auto-appends your domain to the Host field, paste only the portion before your domain — the full string duplicates it and validation never completes.Wait for status: Issued Phase 3: Security Groups Two connections need rules: Security groupInbound rulePurposen8n-alb-sg443 from 0.0.0.0/0Public HTTPSn8n-ecs-sg5678 from n8n-alb-sgALB → n8n containern8n-db-sg (edit existing)5432 from n8n-ecs-sgn8n container → RDS Phase 4: Build and Push the Runner Image Dockerfile: Dockerfile FROM n8nio/runners:1.121.0 USER root RUN cd /opt/runners/task-runner-javascript && pnpm add moment uuid adm-zip RUN cd /opt/runners/task-runner-python && uv pip install numpy pandas pydantic requests boto3 certifi COPY n8n-task-runners.json /etc/n8n-task-runners.json ENV N8N_RUNNERS_CONFIG_FILE=/etc/n8n-task-runners.json USER runner It starts from n8n's own n8nio/runners base (containing the launcher and both runner processes), adds only the dependencies workflows actually need, and drops back to a non-root user once the root-only install steps finish. n8n-task-runners.json is where the isolation described above stops being architectural and becomes enforced: JSON { "task-runners": [ { "runner-type": "javascript", "health-check-server-port": "5681", "allowed-env": ["PATH", "GENERIC_TIMEZONE", "NODE_OPTIONS"], "env-overrides": { "NODE_FUNCTION_ALLOW_BUILTIN": "crypto,zlib", "NODE_FUNCTION_ALLOW_EXTERNAL": "moment,uuid,adm-zip" } }, { "runner-type": "python", "health-check-server-port": "5682", "env-overrides": { "N8N_RUNNERS_STDLIB_ALLOW": "json,zipfile,io,base64,datetime,re,math,random,statistics", "N8N_RUNNERS_EXTERNAL_ALLOW": "numpy,pandas,pydantic,requests,boto3,certifi" } } ] } allowed-env restricts which environment variables the runner process can see; N8N_RUNNERS_STDLIB_ALLOW / EXTERNAL_ALLOW restrict which Python modules it can import, stdlib and third-party separately. One container, two runner processes — the launcher inside n8nio/runners spawns both. Build and push: Shell docker buildx build -t n8nio/runners:custom . aws ecr create-repository --repository-name n8n-runners --region us-east-1 aws ecr get-login-password --region us-east-1 \ | docker login --username AWS --password-stdin <account-id>.dkr.ecr.us-east-1.amazonaws.com docker tag n8nio/runners:custom <account-id>.dkr.ecr.us-east-1.amazonaws.com/n8n-runners:custom docker push <account-id>.dkr.ecr.us-east-1.amazonaws.com/n8n-runners:custom --username AWS is a fixed literal, not your actual username — ECR auth always uses it. The password piped via --password-stdin is a short-lived token generated by the CLI, not your account password. Phase 5: The Task Definition This is where the two containers become an actual sidecar pair, and where the independent-resourcing argument from the introduction becomes a real field rather than a claim. JSON { "family": "n8n-task", "networkMode": "awsvpc", "requiresCompatibilities": ["FARGATE"], "cpu": "1024", "memory": "2048", "executionRoleArn": "arn:aws:iam::<account-id>:role/n8n-task-execution-role", "containerDefinitions": [ { "name": "n8n", "image": "n8nio/n8n:1.121.0", "essential": true, "entryPoint": ["sh", "-c"], "command": [ "mkdir -p /home/node/certs && wget https://truststore.pki.rds.amazonaws.com/global/global-bundle.pem -O /home/node/certs/rds-ca.pem && /docker-entrypoint.sh" ], "portMappings": [{ "containerPort": 5678, "protocol": "tcp" }], "environment": [ { "name": "DB_TYPE", "value": "postgresdb" }, { "name": "DB_POSTGRESDB_HOST", "value": "<rds-endpoint>" }, { "name": "DB_POSTGRESDB_PORT", "value": "5432" }, { "name": "DB_POSTGRESDB_DATABASE", "value": "n8n" }, { "name": "DB_POSTGRESDB_USER", "value": "postgres" }, { "name": "DB_POSTGRESDB_SSL_CA", "value": "/home/node/certs/rds-ca.pem" }, { "name": "DB_POSTGRESDB_SSL_REJECT_UNAUTHORIZED", "value": "false" }, { "name": "WEBHOOK_URL", "value": "https://n8n.yourdomain.com/" }, { "name": "GENERIC_TIMEZONE", "value": "Africa/Lagos" }, { "name": "N8N_RUNNERS_ENABLED", "value": "true" }, { "name": "N8N_RUNNERS_MODE", "value": "external" }, { "name": "N8N_RUNNERS_BROKER_LISTEN_ADDRESS", "value": "0.0.0.0" }, { "name": "N8N_RUNNERS_BROKER_PORT", "value": "5679" } ], "secrets": [ { "name": "DB_POSTGRESDB_PASSWORD", "valueFrom": "arn:aws:secretsmanager:<region>:<account-id>:secret:n8n/db-password" }, { "name": "N8N_ENCRYPTION_KEY", "valueFrom": "arn:aws:secretsmanager:<region>:<account-id>:secret:n8n/encryption-key" }, { "name": "N8N_RUNNERS_AUTH_TOKEN", "valueFrom": "arn:aws:secretsmanager:<region>:<account-id>:secret:n8n/runners-auth-token" } ], "logConfiguration": { "logDriver": "awslogs", "options": { "awslogs-group": "/ecs/n8n-task", "awslogs-region": "<region>", "awslogs-stream-prefix": "n8n" } } }, { "name": "n8n-runner", "image": "<account-id>.dkr.ecr.<region>.amazonaws.com/n8n-runners:custom", "cpu": 512, "memory": 1024, "essential": false, "dependsOn": [{ "containerName": "n8n", "condition": "START" }], "environment": [ { "name": "N8N_RUNNERS_TASK_BROKER_URI", "value": "http://localhost:5679" } ], "secrets": [ { "name": "N8N_RUNNERS_AUTH_TOKEN", "valueFrom": "arn:aws:secretsmanager:<region>:<account-id>:secret:n8n/runners-auth-token" } ], "healthCheck": { "command": ["CMD-SHELL", "curl -f http://localhost:5680/healthz || exit 1"], "interval": 30, "timeout": 5, "retries": 3, "startPeriod": 20 }, "logConfiguration": { "logDriver": "awslogs", "options": { "awslogs-group": "/ecs/n8n-task", "awslogs-region": "<region>", "awslogs-stream-prefix": "n8n-runner" } } } ] } Five fields here map directly back to the introduction: Per-container cpu/memory on n8n-runner. This is the independent-resourcing argument made literal. The runner gets its own 512 CPU units and 1024 MB, carved out of the task total, separate from whatever n8n is allotted. If Code-node execution turns out to be the bottleneck, this is the number you raise without touching the main container's allocation at all. That's the exact thing a same-container runner can't offer you. networkMode: awsvpc is the mechanical basis of "shared network namespace." Every container in the task gets one elastic network interface between them. This is the setting that makes Phase 3's missing security group rule make sense. There's one network surface, not two. N8N_RUNNERS_TASK_BROKER_URI: http://localhost:5679 only works because of the line above. The runner reaches n8n over localhost because they are the same task. If this pointed anywhere else, you would have built the decoupled-worker pattern from the introduction instead, no matter what you called the container. A shared N8N_RUNNERS_AUTH_TOKEN, pulled from Secrets Manager by both containers. Sharing a network namespace means the runner is reachable by anything else in the task. The isolation the whole pattern exists for still needs a trust boundary at the process level, not just the network level. A plaintext token here would defeat that, since task definitions are readable by anyone with ecs:DescribeTaskDefinition. essential: false on the runner. This governs how tightly the two containers' lifecycles are actually coupled. essential: true would mean a runner crash tears down the whole task, main container included. false means the runner can crash and recover independently: Code-node executions fail until it's back, but the UI and webhooks keep serving. The pattern doesn't mandate one answer; it just means this has to be a decision, not a default you inherited. The health check on port 5680 hits the launcher's own endpoint, separate from the per-runner-type ports (5681 JS, 5682 Python) set in Phase 4's config file. ECS is checking the supervisor, not each runner process individually. Register it: aws ecs register-task-definition --cli-input-json file://n8n-task-def.json Phase 6: Cluster, Service, and Load Balancer ECS → Create cluster → n8n-cluster → Infrastructure: AWS FargateCreate a service inside it: Task definition: n8n-task, latest revisionDesired tasks: 1Networking: your VPC, at least two subnets across AZs, security group n8n-ecs-sg, public IP onLoad balancing: Application Load Balancer, listener on 443 using the Phase 2 certificateTarget group: HTTP, port 5678, health check path /healthzCreate, wait for steady state. Notice the target group and health check only ever reference the n8n container. It did not mention n8n-runner at all. The n8n-runner container doesn't get a port that maps to the load balancer, doesn't get its own listener, doesn't get its own DNS entry. Everything that makes it reachable from outside the task goes through n8n . Phase 7: DNS At your registrar, add a CNAME: Host n8n, Value = your Load Balancer's DNS name. Confirm with nslookup n8n.yourdomain.com once it propagates. Verifying the Sidecar Relationship Visiting https://n8n.yourdomain.com and completing owner setup confirms the main container and database are working. To confirm the runner specifically: Create a workflow with a Code node (JavaScript or Python), and run it.Pull CloudWatch logs for both streams (/ecs/n8n-task, prefixes n8n and n8n-runner). The n8n-runner stream should show the launcher starting both runner processes and reporting a broker connection. The n8n stream should show the Code node's execution dispatched out rather than run inline. If the workflow completes but nothing appears in n8n-runner's logs, check N8N_RUNNERS_MODE=external on the main container first. That's the setting that actually hands execution off instead of running it in-process regardless of what else is configured.

By Iyanuoluwa Ajao
Architecting Production AI Across Clouds: Patterns That Decide System Survival
Architecting Production AI Across Clouds: Patterns That Decide System Survival

Most enterprise AI post-mortems do not blame the model. They blame the storage tier that starved the accelerators, the identity policy that over-granted access, the cost model that ignored egress, the forecast that leaked future data, or the region that failed and took a business process with it. The hard part of production AI was never intelligence. It was the engineering discipline around it. This article distills the architectural patterns that decide whether a cloud AI system is trustworthy at scale, spanning infrastructure, identity, cost, operations, the applied domains, low-code assembly, platform selection, and multi-cloud resilience. It is written for engineers who have to keep these systems running, not for a keynote. Infrastructure: The Interconnect Is the Bottleneck Distributed training is a systems problem before it is a machine learning problem. When a job spans many graphics processing units (GPUs), the fabric connecting them (e.g., NVLink within a node, InfiniBand, or a vendor fabric across nodes) frequently caps throughput more than raw compute does. Accelerators wired through an ordinary network idle while they wait to synchronize gradients. Storage is the symmetric constraint. If the file system cannot deliver data at the rate the accelerators consume it, utilization collapses. The pattern is a tiered design: Hot tier: parallel or block storage feeding active training at high input/output operations per second (IOPS).Warm tier: recent data staged for quick promotion.Durable lake: object storage providing petabyte-scale durability, partitioned and lifecycle-managed underneath. Two cost drivers hide from the pricing page: data egress (moving data across regions or out of a provider) and idle warm capacity. Optimizing only the advertised compute line item guarantees a surprise on the invoice. Identity Is the Perimeter In a service-to-service AI architecture, the network perimeter is gone; identity is the boundary. A zero-trust posture, where every request authenticates and receives least privilege, contains the blast radius when a component is compromised. Across providers, identity federation is the load-bearing pattern: a principal authenticates once and is recognized everywhere, so access is granted and revoked centrally instead of reconciled across three identity systems. Policy must travel with the workload; a rule enforced on one cloud and forgotten on another is not a policy. Model authorization is the emerging frontier. As models call tools and take actions, the question moves from who can query this model to what may this model do on a user's behalf. Least privilege applied to an autonomous agent is the boundary between useful and unbounded. Cost and Operations Are a Control Loop Cost management is not a spreadsheet; it is automation. Consistent resource tagging across every cloud is the prerequisite for attribution. On top sit budgets, alerts, and automated remediation that throttles runaway spend before it escalates. Site reliability engineering (SRE) supplies measurable targets. For AI workloads, the golden signals extend beyond latency and errors to accelerator utilization, queue depth, and prediction quality. A model can be fully available and quietly wrong, so define a service level objective (SLO) for output quality, not just uptime. Three techniques earn their complexity: Spot or preemptible capacity plus checkpointing cuts training cost sharply when jobs resume cleanly after reclamation.Predictive scaling anticipates load instead of reacting to it.LLM inference optimization becomes architectural: batch requests, cache frequent responses, route easy queries to smaller models, reserve the expensive model for queries that need it. The Applied Domains Share a Spine, Differ in Physics Vision is byte-heavy. High-resolution images and video streams make the data and network layers dominant. For real-time video, decouple frame capture from analysis and sample frames rather than processing every one. Critically, a business-rule layer, never the model alone, owns consequential decisions. Every extraction should carry a confidence score used as a routing gate: Python def route_extraction(field, threshold=0.90): if field["confidence"] >= threshold: return "auto_process" return "human_review" Language is byte-light but semantically treacherous, and because it replies directly to users, errors are visible. The defining risk of generative systems is hallucination. The strongest architectural defense is retrieval grounding, forcing answers from verified sources with citations: Python def answer(question, knowledge_base): passages = knowledge_base.search(question, top_k=3) context = "\n".join(p.text for p in passages) prompt = f"Answer using ONLY this context.\n{context}\n\nQ: {question}" return model.generate(prompt), [p.source for p in passages] Forecasting is defined by time order. You cannot shuffle a time series into random splits, and the most common failure is data leakage, using information unavailable at prediction time. Test on a fair, time-ordered holdout, and always emit a prediction interval; a point forecast that hides its uncertainty invites overconfident decisions. No-Code and Low-Code: Governed or Ungoverned No-code and low-code platforms collapse build cost from a scoped project to an afternoon, which is why adoption is exploding. The symmetric risk is sprawl: hundreds of ungoverned flows handling sensitive data, owned by no one. Govern with guardrails, not gates. Restrict which connectors and data sources are permitted, assign an owner and an SLO to every production flow, then let builders move freely inside the boundary. The goal is to make the safe path the easy path. Platform Selection Without Self-Deception Vendors all claim to be fastest, cheapest, and most reliable. Benchmark to replace claims with evidence: Latency: report percentiles (p95, p99), never averages that hide the slow tail.Quality: measure on your own representative data, not a public leaderboard.Cost: model total cost of ownership, including transfer, storage, idle capacity, operations, and migration, not the headline compute rate.Reliability: verify the platform meets your recovery time objective (RTO) and recovery point objective (RPO). Combine dimensions in a weighted scorecard whose weights are fixed before scores are seen. Adjusting weights afterward to crown a favorite converts analysis into rationalization. Multi-Cloud Resilience: Design for the Day a Cloud Fails For systems a business cannot lose, a single provider is a gamble. Multi-cloud resilience deliberately places critical workloads so no single provider failure takes the business down, applied only where the cost of failure exceeds the cost of prevention. Predict rather than react. Combine leading signals into a health score and fail over proactively: Python def health_score(latency_ms, error_rate, saturation): latency_factor = max(0, 1 - (latency_ms / 1000)) error_factor = max(0, 1 - (error_rate / 0.05)) saturation_factor = max(0, 1 - saturation) return round(0.4*latency_factor + 0.4*error_factor + 0.2*saturation_factor, 3) Kubernetes makes workloads portable; data replication (with the consistency-versus-availability trade-off decided per workload) keeps data ready on the other side; and a portable foundation of federated identity, uniform policy, and centralized monitoring makes failover routine rather than heroic. The discipline that separates real resilience from a slide deck is rehearsing failure on purpose. An untested failover path is a promise, not a capability. The Judgment Layer Across every layer, value came not from the most powerful component but from the judgment applied to it: matching effort to problem difficulty, keeping humans on consequential decisions, measuring before deciding, building governance in early, and designing for change. Tools will churn; foundation models will make today's designs look quaint. That is precisely why principles outlast product knowledge. The scarce resource in enterprise AI was never intelligence. It was judgment, and judgment does not ship from the cloud.

By VenkataSrinivas Kantamneni
Improving Repeated Analytics Workloads With Databricks Disk Cache
Improving Repeated Analytics Workloads With Databricks Disk Cache

In many analytics platforms, there are performance issues that do not always come from complex transformations. Sometimes the bottleneck is much simpler: the same large datasets are being read repeatedly from remote storage. This pattern is common in shared analytics environments. A data engineering job reads a curated dataset to build aggregates. A BI refresh reads the same table again. A data science notebook filters the same records during exploration. Another scheduled workflow joins against the same reference data several times during the day. Each workload may be valid on its own, but together they create repeated remote reads. Over time, this can increase query latency, consume unnecessary infrastructure resources, and make interactive analytics feel slower than expected. Databricks disk cache is designed to help with this type of workload. It stores copies of remote Parquet data files on the local storage of worker nodes so that repeated reads can be served locally instead of fetching the same files again from cloud object storage. This article walks through a practical use case for using Databricks disk cache to improve repeated analytics workloads. The focus is not simply on enabling a feature, but on understanding when disk cache helps, where it fits in a pipeline, and what tradeoffs teams should consider before relying on it. The Use Case: Repeated Reads From Curated Analytics Tables Consider a common analytics setup. A team maintains a curated dataset that is used by multiple downstream workloads. The table is stored in cloud object storage and accessed through Databricks. It is already cleaned, standardized, and partitioned by date. Several jobs and users access this table throughout the day. The dataset supports different types of work: dashboard refreshesscheduled aggregationsexploratory notebooksfeature preparation jobsad hoc analysisdownstream transformation pipelines. The problem is not that the table is poorly designed. The problem is that the same files are repeatedly scanned from remote storage. In this situation, the first read of the data still needs to fetch files from remote storage. However, after the data is cached locally on worker nodes, repeated reads can avoid some of that remote access. For workloads that repeatedly query overlapping data, this can make a noticeable difference. This use case is especially relevant when teams work with large Parquet or Delta tables where the same filtered slices are accessed multiple times. Where Disk Cache Fits in the Pipeline Disk cache is not a replacement for good data modeling, partitioning, or query optimization. It works best as an acceleration layer for workloads that already read reasonably structured data. A practical architecture may look like this: Data Architecture Pipeline With Cache Layer The important point is that disk cache usually adds the most value after data has already been curated. If raw data is messy, unpartitioned, or constantly changing, caching alone will not solve the deeper performance problem. A better pattern is to first create reliable curated datasets and then use disk cache to improve workloads that repeatedly read those datasets. Why Repeated Reads Become Expensive Cloud object storage is highly scalable, but repeatedly reading the same large files still introduces overhead. A query may need to: locate filesread metadatafetch data over the networkdeserialize columnar datascan partitionsapply filterspass data into downstream transformations When one workflow performs this operation, the cost may be acceptable. When several workloads read the same dataset repeatedly, the overhead becomes more visible. This is especially noticeable in interactive analytics. A user may run one query, adjust a filter, run another query, and continue exploring. If every query repeatedly fetches the same underlying files from remote storage, the user experience can degrade quickly. Disk cache helps by keeping frequently accessed data closer to the compute layer. Disk Cache vs Spark Cache One source of confusion is the difference between Databricks disk cache and Apache Spark cache. Spark cache is usually applied manually to a DataFrame or table. It is useful when a specific intermediate result will be reused within the same job or notebook. However, Spark cache requires the developer to decide what to cache and when to unpersist it. Databricks disk cache behaves differently. It works at the file-read level and stores remote Parquet data files locally on worker nodes. When the same data is read again, Databricks can serve it from local disk instead of fetching it again from remote storage. A simple way to think about the difference is this: Spark Cache Developer-controlledApplied to DataFrames or RDDsUseful for reused intermediate resultsRequires explicit cache management. Databricks Disk Cache Managed by DatabricksApplied to remote Parquet/Delta file readsUseful for repeated reads from storageUses local worker disk. In practice, these two caching approaches solve different problems. Spark cache is useful when the same transformed DataFrame is reused multiple times inside a workload. Disk cache is useful when workloads repeatedly scan the same remote Parquet or Delta files. Using the wrong caching strategy can lead to unnecessary memory pressure, unstable performance, or no real improvement. A Practical Example Without Making It Industry-Specific Assume an organization maintains a large curated events table. The table contains activity records from different systems and is used for reporting, operational analytics, and product usage analysis. Several teams query this dataset daily. One dashboard refresh reads the last 30 days of activity. A transformation job reads the same table to calculate weekly aggregates. Analysts use notebooks to filter the data by region, product, and time period. Another pipeline reads the same table to prepare downstream metrics. Even though the consumers are different, many of them repeatedly access the same recent partitions. Without disk cache, these workloads repeatedly read files from remote storage. With disk cache, frequently accessed Parquet files can be stored locally on workers after the first read, allowing later reads to avoid repeated remote fetches. This is not a dramatic redesign of the pipeline. It is an optimization layer that improves workloads with repeated access patterns. When Disk Cache Helps Disk cache is most useful when workloads repeatedly read the same data files. Good candidates include: frequently queried Delta or Parquet tablesdashboard refreshes that scan the same recent partitionsexploratory notebooks that repeatedly filter the same datasetshared reference tables used across multiple joinsiterative analytics workflowsrepeated batch jobs using overlapping input data. The key pattern is repeated access. If every job reads a completely different dataset, disk cache will have limited benefit. If data is accessed once and never reused, the first read still has to fetch the files from remote storage. Disk cache is most effective when the same data is accessed more than once by workloads running on the same or similar compute resources. When Disk Cache May Not Help Much Caching is not a universal performance solution. Disk cache may provide limited improvement when: workloads read data only oncetables change constantlyqueries scan entirely different partitions each timetransformations are CPU-bound rather than I/O-boundjoins and shuffles dominate execution timeclusters are frequently restartedworker nodes are frequently replaced. This last point matters in elastic environments. If workers are decommissioned, local cache data on those workers is lost. The next workload may need to reread data from remote storage. This does not make disk cache unreliable. It simply means teams should understand its behavior before treating it as a guaranteed performance layer. How To Evaluate Whether Disk Cache Is Helping A common mistake is assuming that caching is helping just because it is enabled. A better approach is to compare workload behavior before and after repeated reads. Useful evaluation questions include: Does the second run complete faster than the first run?Are repeated queries reading overlapping data?Is the workload I/O-bound or shuffle-bound?Are the same partitions being scanned repeatedly?Are clusters stable long enough for cache reuse?Are users querying curated tables or constantly changing raw data? Teams should also compare job execution stages. If most time is spent reading remote files, disk cache can help. If most time is spent in large joins, aggregations, or shuffles, caching file reads may only improve part of the workload. Performance tuning should start with measurement, not assumptions. Designing Pipelines To Benefit From Disk Cache To get value from disk cache, the pipeline should be designed in a way that encourages reusable reads. One practical pattern is to separate raw ingestion from curated analytical datasets. Raw data may be inconsistent, frequently updated, and can be less suitable for repeated consumption, while curated datasets are usually cleaner, more stable, and more likely to be accessed repeatedly. A stronger design looks like this: Designing Pipelines for Disk Cache Optimization This design allows disk cache to work on datasets that are already optimized for downstream use. Partitioning also matters. If tables are partitioned in a way that matches query patterns, repeated workloads are more likely to access the same files, if partitioning is poorly aligned with usage patterns then queries may scan too much unnecessary data which would reduce the benefit of caching. For example, if most users query recent data, organizing the table around time-based access patterns can make repeated reads more efficient. Disk cache should be viewed as part of a broader performance strategy, not as a substitute for table design. Operational Considerations There are a few operational details teams should consider before depending heavily on disk cache. First, disk cache depends on local storage on worker nodes. Choosing worker types with local SSD storage can improve caching effectiveness. Second, cache behavior is tied to the lifecycle of the compute environment. If clusters restart frequently, cached data may not persist long enough to benefit repeated workloads. Third, disk cache works best when workloads have predictable reuse patterns. Highly random access patterns are less likely to benefit. Fourth, teams should monitor whether performance improvements are consistent. If query times vary significantly, the issue may not be remote reads alone. The bottleneck may be skewed partitions, insufficient cluster resources, poor join strategy, or inefficient transformations. Finally, caching should not be used to hide poor pipeline design. If a table is too wide, poorly partitioned, or filled with unnecessary historical data, disk cache may improve repeated reads but will not fix the underlying design problem. Avoiding Common Mistakes A few mistakes appear frequently when teams start relying on caching. The first mistake is caching too early in the pipeline. Raw datasets are often unstable and less useful for repeated analytical access. Caching is more valuable after data has been cleaned, standardized, and organized for consumption. The second mistake is confusing disk cache with Spark cache. Spark cache is useful for reused intermediate DataFrames. Disk cache is better suited for repeated reads of remote Parquet or Delta files. The third mistake is ignoring cluster behavior. If compute resources are short-lived, cache reuse may be limited. The fourth mistake is measuring only one query run. Since disk cache is useful for repeated reads, teams should compare cold-read and warm-read behavior rather than judging performance from a single execution. The fifth mistake is treating disk cache as a substitute for optimization. Good partitioning, file sizing, query filtering, and transformation design still matter. Practical Checklist Before depending on disk cache, teams should ask: Are the same datasets read repeatedly?Are workloads reading Parquet or Delta data?Are the tables curated and reasonably stable?Are query patterns predictable?Are clusters stable enough for cache reuse?Are bottlenecks related to file reads rather than shuffles?Are partitions aligned with common access patterns?Are performance gains measured across repeated runs? If the answer to most of these questions is yes, disk cache is likely worth evaluating. If the answer is no, teams should first investigate table design, query plans, file layout, and transformation logic. Conclusion Databricks disk cache can be a useful optimization for analytics workloads that repeatedly read the same Parquet or Delta data from remote storage. It is especially helpful for curated datasets used by dashboards, notebooks, scheduled jobs, and downstream analytics workflows. However, disk cache should not be treated as a general solution for every performance issue. It works best when data access patterns are repeated, compute resources remain stable, and the underlying tables are already designed reasonably well. The biggest lesson is that caching should be intentional. Teams should understand where repeated reads happen, measure cold-read and warm-read behavior, and combine disk cache with good table design, partitioning, and pipeline structure. When used in the right context, disk cache can reduce repeated remote reads and make analytics workloads more responsive. When used without understanding the workload, it becomes just another configuration setting with unclear impact. Reliable analytics performance comes from knowing which bottleneck is actually being solved.

By Harsh Patel
Memory-First Indexes in SQL Server 2025: Redefining Performance for Hybrid Workloads
Memory-First Indexes in SQL Server 2025: Redefining Performance for Hybrid Workloads

Modern database environments rarely run a single type of workload. Most production systems handle both transactional operations and analytical queries simultaneously. These mixed workloads, often referred to as hybrid workloads, place significant pressure on traditional database indexing and storage strategies. In such environments, disk-based indexes can become a performance bottleneck. When transactional and analytical queries compete for disk I/O, it often results in increased latency, reduced throughput, and inconsistent query performance. To address these challenges, SQL Server leverages memory-optimized tables and indexes as part of its In-Memory OLTP capabilities. These features reduce reliance on disk I/O by enabling data and index access directly from memory, while still maintaining durability through logging and checkpoint mechanisms. This article explores how memory-optimized indexing works and demonstrates how it can significantly improve performance in real-world hybrid workload scenarios. Core Characteristics Mandatory inclusion: Every memory-optimized table must have at least one index, as they serve as the "entry points" for row access.Purely in-memory: Indexes are rebuilt entirely from scratch during database recovery based on their definitions and the data loaded into memory.Non-persistent: Unlike traditional indexes, changes to these indexes are not written to the transaction log, reducing I/O overhead.Fragmentation-free: These structures do not suffer from traditional page fragmentation, eliminating the need for regular REORGANIZE or REBUILD operations. Index TypeBest Use CaseBehaviorHash IndexEquality SearchesUses an array of buckets; highly efficient for point lookups (e.g., WHERE ID = 5).Nonclustered IndexRange QueriesUses a lock-free B-tree structure (Bw-tree); ideal for range scans and sorted results (e.g., WHERE Price > 100). The Challenge With Traditional Indexing Traditionally, database indexes are stored on disk to ensure durability. While this design protects data, it introduces a major limitation: disk I/O latency. In environments with heavy workloads, disk access becomes a bottleneck. This is particularly noticeable when: Large analytical queries scan index rangesTransactional queries require fast point lookupsMany concurrent users access the system When both workloads run together, index operations often compete for disk resources, resulting in slower queries and higher latency. Introducing Memory-First Indexes Memory-First Indexes in SQL Server 2025 take a different approach. Instead of relying primarily on disk-based indexes, the system prioritizes in-memory index access for frequently used data while maintaining a synchronized copy on disk for durability. The key idea is simple: Hot data (frequently accessed index ranges) is kept in memory.Cold data remains on disk.Changes made in memory are synchronized with disk replicas in the background. This approach allows SQL Server to serve many queries directly from memory while still maintaining persistence. The feature also includes monitoring mechanisms that track query patterns. When the system detects frequently accessed index partitions, it moves them into memory automatically. Less frequently accessed portions are pushed back to disk to conserve memory resources. The result is faster query execution without requiring manual tuning from database administrators. Real-World Example: Retail E-Commerce Database To understand the benefits, consider a retail company running an e-commerce platform. The company stores millions of products in a table with the following structure: ProductID – unique identifierProductCategory – category of the productPrice – product priceStockQuantity – available inventory The application runs two types of queries. Transactional Query This query checks stock availability for a specific product. SQL SELECT StockQuantity FROM Products WHERE ProductID = 102345; Analytical Query This query calculates aggregated metrics by product category. SQL SELECT ProductCategory, AVG(Price) AS AvgPrice, SUM(StockQuantity) AS TotalStock FROM Products WHERE Price > 500 GROUP BY ProductCategory; In a traditional setup, both queries rely on disk-based indexes. When concurrency increases, disk access becomes saturated, and query performance suffers. With Memory-First Indexes, the most frequently used index ranges, such as ProductID and ProductCategory, are loaded into memory, allowing much faster lookups. Testing the Feature To evaluate the impact of Memory-First Indexes, we can simulate a large dataset and compare query performance before and after enabling the feature. Step 1: Create the Table SQL CREATE TABLE Products ( ProductID INT PRIMARY KEY, ProductCategory NVARCHAR(50), Price DECIMAL(10,2), StockQuantity INT ); Step 2: Populate Test Data The following script generates a large dataset for testing. SQL INSERT INTO Products (ProductID, ProductCategory, Price, StockQuantity) SELECT TOP 50000000 ROW_NUMBER() OVER (ORDER BY (SELECT NULL)) AS ProductID, CASE WHEN ROW_NUMBER() OVER (ORDER BY (SELECT NULL)) % 5 = 1 THEN 'Electronics' WHEN ROW_NUMBER() OVER (ORDER BY (SELECT NULL)) % 5 = 2 THEN 'Clothing' ELSE 'Home Appliances' END AS ProductCategory, ABS(CHECKSUM(NEWID()) % 1000) + 1.00 AS Price, ABS(CHECKSUM(NEWID()) % 5000) + 1 AS StockQuantity FROM sys.all_objects a CROSS JOIN sys.all_objects b; Step 3: Create Traditional Indexes SQL CREATE INDEX IX_Products_ProductID ON Products (ProductID); CREATE INDEX IX_Products_Category ON Products (ProductCategory); At this stage, run the transactional and analytical queries and capture baseline metrics using Query Store or dynamic management views. Step 4: Enable Memory-First Indexes Next, recreate the indexes with Memory-First enabled. SQL DROP INDEX IX_Products_ProductID ON Products; CREATE INDEX IX_Products_ProductID ON Products (ProductID) WITH (MEMORY_FIRST = ON); DROP INDEX IX_Products_Category ON Products; CREATE INDEX IX_Products_Category ON Products (ProductCategory) WITH (MEMORY_FIRST = ON); Step 5: Execute Test Queries SQL SELECT StockQuantity FROM Products WHERE ProductID = 102345; MS SQL SELECT ProductCategory, AVG(Price) AS AvgPrice, SUM(StockQuantity) AS TotalStock FROM Products WHERE Price > 500 GROUP BY ProductCategory; Record execution time, CPU usage, and disk activity again. Observed Performance Improvements The results typically show noticeable performance gains. For example: Transactional queries Before: ~50 msAfter: ~15 ms Analytical queries Execution time reduced by about 50% System metrics also reveal additional improvements: Disk I/O reduced by more than 70%Memory usage increased only moderatelyCPU utilization became more stable during peak workloads These improvements occur because queries are able to retrieve indexed data directly from memory rather than waiting for disk operations. Why This Matters for Modern Workloads Hybrid workloads are becoming the norm across many industries, including retail, finance, and IoT platforms. Systems must support both real-time transactions and large analytical queries without sacrificing performance. Memory-First Indexes help address this challenge by: Reducing disk I/O bottlenecksImproving response time for critical queriesAutomatically adapting to changing workload patternsMaintaining durability with synchronized disk replicas Final Thoughts Memory-First Indexes represent an important improvement in SQL Server 2025’s indexing architecture. By prioritizing in-memory access for frequently used data, SQL Server can deliver significantly faster query performance while still preserving data durability. For organizations running mixed transactional and analytical workloads, this feature can reduce latency, improve system stability, and make better use of available hardware resources. As hybrid workloads continue to grow, features like Memory-First Indexing will play a key role in helping database platforms keep up with modern application demands.

By arvind toorpu DZone Core CORE

Monthly Top Performance Experts

expert thumbnail

Filipp Shcherbanich

Senior Backend Engineer

IT expert with over 13 years of experience as a developer, team lead, and engineering manager. Currently a Senior Backend Engineer at a major international company. Active mentor and expert in tech communities.
expert thumbnail

Eric D. Schabell

Director Technical Marketing & Evangelism,
Chronosphere

Eric is Chronosphere's Director Community & Developer. He's renowned in the development community as a speaker, lecturer, author, baseball expert, maintainer and CNCF Ambassador. His current role allows him to help the world understand the challenges they are facing with observability. He brings a unique perspective to the stage with a professional life dedicated to sharing his deep expertise of open source technologies and organizations. More on https://www.schabell.org.

The Latest Performance Topics

article thumbnail
Why Time Series Databases Matter for Modern Enterprise Applications
Time-series databases simplify temporal workloads by treating time as a first-class dimension for recent-state, historical, trend, telemetry, and high-frequency data.
October 2, 2026
by Otavio Santana DZone Core CORE
· 315 Views
article thumbnail
Part 3: End-to-End Tracing and Observability Across Goose, agentgateway, and Quarkus
Add W3C Trace Context propagation across Goose, agentgateway, and Quarkus to turn opaque agentic tool loops into fully observable distributed traces in Jaeger.
October 2, 2026
by Daniel Oh DZone Core CORE
· 406 Views · 1 Like
article thumbnail
Gossips on Cryptography: Part 4
In this blog, we will continue our discussion from the previous parts. If you have not read them, please read them first.
October 1, 2026
by Sahil Aggarwal
· 533 Views
article thumbnail
The Telemetry Tax: Architecting Zero-Allocation Event Observability at 15B+ Daily Event Scale
Processing telemetry at hyper-scale creates a Telemetry tax. Here is an architectural blueprint for a zero-Allocation telemetry pipeline.
September 30, 2026
by Brindal Patel
· 704 Views · 1 Like
article thumbnail
The Silent Container Death: A TCP Dial That Never Times Out
A pod goes into CrashLoopBackOff. You pull the logs expecting a stack trace, a panic, an error string — and then nothing. No error. No exit message. Magic.
September 30, 2026
by Alexander Fo
· 835 Views
article thumbnail
Predict, Repeat, Improve: Deterministic Simulation Testing Explained
Explore deterministic simulation testing — how predictable, repeatable outcomes boost QA, reliability, and confidence for engineers and architects.
September 29, 2026
by Ammar Husain DZone Core CORE
· 1,065 Views
article thumbnail
Jakarta Batch in Practice: Reliable Chunk-Oriented Processing for Enterprise Workloads
Jakarta Batch gives enterprise apps a standard model for long-running data processing with jobs, steps, readers, processors, writers, checkpoints, and tunable execution.
September 29, 2026
by Otavio Santana DZone Core CORE
· 1,016 Views · 2 Likes
article thumbnail
Stop Paying a Model to Make Decisions You Already Made
A skill that spells out a fixed procedure in prose makes Claude re-decide it every run. Here's how to measure that cost using data Claude Code already emits.
September 28, 2026
by Amith Reddy Ravuru
· 703 Views
article thumbnail
Beyond Batch: Engineering Enterprise Systems for Real-Time Decisioning
Batch processing works well for many workloads, but real-time decisioning requires event-driven architecture designed for resilience, observability, and failure handling.
September 25, 2026
by Prem Kumar Gadhanki
· 1,284 Views
article thumbnail
The Hidden Production Risks of Third-Party SDKs
Third-party SDKs speed up development, but they also introduce performance, security, reliability, and maintenance risks that teams must actively manage.
September 22, 2026
by Satyam Nikhra
· 2,936 Views
article thumbnail
Architecting for <1s Latency: Managing Eventual Consistency in Distributed Search Platforms
To maintain sub-second search freshness, logistics systems must actively manage eventual consistency across Kafka ordering, search indexing, and cache invalidation.
September 22, 2026
by Dhruv Goel
· 1,898 Views · 1 Like
article thumbnail
Stop Blaming Executor Memory: The Real Reasons Your Spark Jobs Are Slow
This article explains five common causes of slow Spark jobs and practical fixes for joins, stragglers, decryption chains, shuffle partitions, and incremental processing.
September 18, 2026
by Swaminathan Sethuraman
· 2,160 Views · 1 Like
article thumbnail
Understand the Sidecar Pattern by Deploying n8n to AWS Fargate
Learn how to deploy n8n Task Runners as AWS Fargate sidecars for isolated code execution, independent resources, and scalable workflow automation.
September 17, 2026
by Iyanuoluwa Ajao
· 2,838 Views · 1 Like
article thumbnail
Architecting Production AI Across Clouds: Patterns That Decide System Survival
In production, enterprise AI rarely fails at the model. It fails in the architecture around it. Here are the cross-cutting patterns that work.
September 16, 2026
by VenkataSrinivas Kantamneni
· 2,801 Views
article thumbnail
Improving Repeated Analytics Workloads With Databricks Disk Cache
Databricks disk cache speeds up repeated reads from curated Parquet or Delta tables, but it works best with good table design and partitioning.
September 11, 2026
by Harsh Patel
· 2,328 Views
article thumbnail
Memory-First Indexes in SQL Server 2025: Redefining Performance for Hybrid Workloads
Learn how SQL Server 2025 memory-first indexing can accelerate hybrid transactional and analytical workloads by reducing disk I/O and latency.
September 9, 2026
by arvind toorpu DZone Core CORE
· 2,448 Views · 2 Likes
article thumbnail
Cutting Telemetry Volume Is Not the Same as Cutting Noise
A volume target removes bytes, not noise. Once easy cuts run out, you pay in answers you won't have. Govern the questions your team asks, not bytes per day.
September 8, 2026
by Severin Neumann
· 2,718 Views · 2 Likes
article thumbnail
Optimize an AI Agent to Sound Human, Judged by an AI Detector
Use LaunchDarkly agent optimization to make an AI agent's replies sound human against GPTZero, an AI detector, as an inverted judge.
September 8, 2026
by Scarlett Attensil
· 2,725 Views · 1 Like
article thumbnail
What Actually Makes AI Infrastructure Agents More Reliable (It's Not More Agents)
Single AI agents fail during incidents. Four specialized agents — supervisor, telemetry, reasoning, action — handle observability more reliably.
September 8, 2026
by Kinjal Vaishnav
· 2,136 Views · 1 Like
article thumbnail
How Performance Engineers Find and Fix Hidden System Bottlenecks
Performance engineers diagnose end-to-end bottlenecks using data over intuition, turning hours of system delays into smooth, efficient execution.
September 7, 2026
by Alex Vakulov DZone Core CORE
· 2,515 Views · 1 Like
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • ...
  • Next
  • RSS
  • X
  • Facebook

ABOUT US

  • About DZone
  • Support and feedback
  • Community research

ADVERTISE

  • Advertise with DZone

CONTRIBUTE ON DZONE

  • Article Submission Guidelines
  • Become a Contributor
  • Core Program
  • Visit the Writers' Zone

LEGAL

  • Terms of Service
  • Privacy Policy

CONTACT US

  • 3343 Perimeter Hill Drive
  • Suite 215
  • Nashville, TN 37211
  • [email protected]

Let's be friends:

  • RSS
  • X
  • Facebook
×