Big data comprises datasets that are massive, varied, complex, and can't be handled traditionally. Big data can include both structured and unstructured data, and it is often stored in data lakes or data warehouses. As organizations grow, big data becomes increasingly more crucial for gathering business insights and analytics. The Big Data Zone contains the resources you need for understanding data storage, data modeling, ELT, ETL, and more.
Why Databricks and Snowflake Speak the Kafka Protocol: Ingestion vs Architecture
OpenAI ‘o’ Leak: What We Know About ChatGPT’s Always-On Assistant Before DevDay
The modern enterprise generates and consumes unprecedented volumes of data across operational systems, customer interactions, partner ecosystems, cloud applications, IoT devices, and AI platforms. At the same time, AI systems are becoming major consumers of enterprise data, making decisions, generating content, recommending actions, and automating workflows. Poor data quality is no longer just a reporting issue; it is also an AI issue. Inaccurate, incomplete, or poorly governed data can produce biased outcomes, regulatory violations, AI hallucinations, and flawed business decisions. Traditional data governance programs were primarily designed to support business intelligence and regulatory compliance. However, the AI era introduces new requirements around model governance, explainability, lineage, ethical AI, data observability, and autonomous decision-making. Organizations must therefore evolve toward a unified data and AI governance model that ensures data can be trusted not only by humans but also by machines. Poor data governance can result in hallucinating AI systems, biased model outcomes, regulatory violations, security breaches, increased operational costs, customer trust erosion, and incorrect business decisions. Data governance has therefore evolved from a compliance function into a strategic business capability. Enterprises that establish trusted, governed, and accessible data foundations will be better positioned to scale AI initiatives, accelerate innovation, and create sustainable competitive advantages. The basic objectives of data governance are: Enhance the agility of data-informed business decisionsFacilitate seamless knowledge sharing across the enterpriseEliminate ambiguity and foster trust in data assetsIncrease data trust, better decision-making, and faster innovation cyclesImprove compliance posture, reduce data duplication, and increase business agility To fully comprehend these objectives, it is essential to first recognize the critical role that data governance plays within an enterprise's broader data management strategy. This white paper explores the challenges, next-generation data capabilities, modern data architecture, and strategic considerations required to build an AI-ready data foundation. Industry Trends of Data Governance According to Gartner, “Any organization in any industry, especially those with very large amounts of data, can use AI for business value.” According to Statista, by 2027 the global market for big data will be worth $103 billion. According to Gartner, 60% of organizations will fail to realize the value of their AI initiatives due to weak data governance frameworks. By 2028, enterprises will increasingly adopt autonomous, AI‑driven governance systems capable of automated policy enforcement, continuous data quality scoring, and real‑time anomaly detection. Gartner forecasts that AI‑driven automation will reduce manual data stewardship tasks by 40% by 2027. Governance models will shift from centralized to federated and hybrid, ultimately evolving toward autonomous domain‑driven governance. Gartner reports that over 60% of enterprises will adopt federated governance by 2027. The rise of AI‑augmented data mesh as a dominant architecture by 2028 (Thoughtworks). AI Trust, Risk, and Security (AI TRiSM) will become the top governance investment area as organizations confront risks related to hallucinations, bias, and regulatory compliance. Gartner predicts that enterprises implementing AI TRiSM will reduce AI‑related risk incidents by 50% by 2026. With the rapid expansion of IoT, 5G, and edge AI, governance must operate in real time. IDC estimates that 30% of enterprise data will be processed at the edge by 2027. AI platforms will embed governance natively, enabling governed prompt engineering, model access, and data contracts. Gartner predicts that 75% of AI platforms will include built‑in governance controls by 2027. Databricks Mosaic AI, Snowflake Cortex, and Microsoft Azure AI’s Responsible AI Dashboard exemplify this trend. Synthetic data will become a regulated and essential component of AI training. Gartner projects that synthetic data will overshadow real data in AI training by 2030. McKinsey estimates that 50% of AI training datasets will include synthetic data by 2028. Global regulations will mandate transparency, lineage, and automated audits. Gartner states that regulatory pressure will be the top driver of data governance investments through 2030. Data contracts will replace traditional API documentation, enforcing schema, SLAs, lineage, and quality. Gartner predicts that data contracts will reduce integration failures by 40% by 2027. Challenges in Data Governance As enterprises expand into multi-cloud environments and increasingly adopt generative AI, governance challenges continue to multiply. The most common data governance challenges faced by enterprises today are, Data explosion: Data exists across multiple, diverse systems throughout the enterprise. Data is spread across structured data, semi-structured data, unstructured content, streaming data, IoT telemetry, computer vision assets, agent-generated content, and AI-generated outputs. Traditional governance frameworks often lack the scalability and automation required to manage such diversity.Data silos: Data is segmented across various platforms, channels, tools, and business units, making it challenging to access across the enterprise. Most data resides across ERP systems, CRM platforms, legacy applications, cloud-native platforms, data warehouses, data lakes, and SaaS applications. This leads to inefficiency, data duplication, and data inconsistency.Data accuracy, completeness, and timeliness: Ensuring data accuracy, completeness, and timeliness remains a challenge.Data quality: Poor oversight of the quality of data coming into an enterprise, as well as its usage throughout the organization, can lead to poor data quality. Common quality challenges include missing values, duplicate records, outdated information, inconsistent definitions, incomplete lineage, and data drift.Regulatory complexity: Managing regulatory compliance, data security, and data privacy presents significant challenges. Enterprises must comply with GDPR, HIPAA, CCPA, PCI-DSS, the EU AI Act, and industry-specific regulations.Data management: Poor data management strategies can result in an enormous amount of data in a completely unmanageable format.Data leakage: Sensitive business information or customer data may be exposed or leaked, leading to misuse. Unsecured data originating from different data sources can lead to data breaches.AI-specific risks: New AI-era governance concerns include algorithmic bias, explainability requirements, training data provenance, prompt governance, LLM hallucinations, and autonomous agent controls. Next Generation Data Capabilities Governance alone does not create value. Enterprises need enterprise data capabilities that make governance operational while enabling innovation and AI adoption. Modern data ecosystems require intelligent platforms capable of discovering, understanding, protecting, and serving data on a scale. Data processing techniques: Unstructured processing covers entity extraction, concept extraction, sentiment analysis, NLP, ontology, etc. To automate portions of the extraction process, Machine Learning techniques are leveraged. Data intelligent platform: It enables natural language queries, AI-powered recommendations, intelligent search, and context-aware discovery.Data products: Data products provide ownership, accountability, defined SLAs, reusability, and business value measurement.On-demand data services: Provide virtualized access to data across the enterprise through way of composable on-demand data services for both online and offline use. It should provide the ability to query in a federated fashion for both online and offline access.Intelligent metadata management: Digital throws data into enterprise systems at a rate that doesn’t allow SMEs to look at data structures and extract metadata. Automated metadata extraction based on ontology is critical. Modern metadata platforms provide automated discovery, classification, catalog generation, lineage tracking, and semantic enrichment.Data fabric: It provides unified data access, cross-platform integration, federated governance, and policy automation.Data mesh: It enables domain ownership, distributed accountability, product-centric thinking, and decentralized governance. Data observability: It focuses on data health monitoring, pipeline performance, anomaly detection, drift identification, and SLA compliance.Real-time analytics: Multi-channel applications and decision management systems are used to capture interactions for digital processes in real-time scenarios. Data archival: Compliance and performance requirements drive the need for archival of both structured and unstructured data. Principles of Data Governance Architecture principles provide a baseline for decision-making across the enterprise. To guide implementation, enterprise data governance principles are categorized into three strategic domains: Value and ownership, security, privacy and ethics, and architecture and quality. Value and Ownership Data as an asset: Data is an enterprise asset with specific, measurable value to the enterprise and must be managed accordingly.Data is shared: Users have access to the data necessary to perform their duties; therefore, data is shared across enterprise functions and business units. Data stewardship: Governance structure must define the owner and those accountable for data-related decisions that are cross-functional. Define the personnel accountable for leadership activities and assign responsibilities to individual contributors or groups of data handlers. Data trustee: Each data element has an assigned trustee accountable for its quality, lifecycle, and compliance. Security, Privacy & Ethics Principles Data security: Data is protected from unauthorized use and disclosure. Data privacy: Privacy and data protection are considered throughout the entire life cycle of the data. All data sharing will conform to relevant regulatory and business requirementsData integrity: Each party to data must be aware of, and abide by, their responsibilities regarding the provision of source data and the obligation to establish and maintain adequate controls over the use of personal or other sensitive data. Data transparency: Governance decisions, policies, and lineage must be transparently documented and clearly communicated across the enterprise. All data-related decisions must be explained clearly to all personnel how, when, and why they are introduced. Architecture & Quality Principles Common vocabulary and data definitions: Data definitions are consistent across the enterprise and understandable to all users.Fit for purpose: Next-generation information ecosystem needs to have fit-for-purpose tools, as no one technology will satisfy all the workloads and processing techniques - E.g., Text Processing, Data Discovery, Dynamic Data Services, High-Performance Analysis, Streaming Analytics, etc.Data metrics: Critical Data Elements (CDEs) of the Business are managed through a lifecycle-oriented data governance process to ensure data quality, with clear metrics and dashboards. As data will reside in many repositories, integrated metadata lineage and PII protection are important. Key Components of Data Governance In the modern era, data management covers both technical requirements and strategic assets for businesses. Efficient data management Strategies help enterprises make informed decisions, improve customer experiences, and drive innovation. Data governance covers the automation of policies, guidelines, principles, and standards for managing data assets. It ensures data quality, accuracy, and compliance with regulatory requirements, building trust in the data. Data governance must be aligned with EA Governance at the enterprise level to realize the business objectives. Some of the open-source data governance tools are Amundsen, DataHub, Apache Atlas, Magda, Open Metadata, Egeria, and TrueData. These tools offer features like Metadata Management, Data Cataloging, and Collaboration to manage data assets effectively. The major components of data governance are: Data qualityData stewardshipData policies and procedures Data security Metadata management Master data management Data storageData privacy and complianceData metrics The following figure depicts the key components of data governance: Figure 1: Key Components of Data Governance Data Quality It helps ensure the accuracy, completeness, and consistency of data. Data quality management involves identifying and correcting errors, standardizing formats, and maintaining a high level of data integrity. Some of the top open-source data quality tools are: Cucumber, Deequ, dbt Core, MobyDQ, Great Expectations, and Soda Core. These tools help automate data validation, data cleaning, and monitoring. Data Stewardship It is about assigning roles and responsibilities related to data management. Data stewards are designated individuals or teams entrusted with overseeing the appropriate use, integrity, and secure storage of enterprise data. They serve as a vital bridge between IT and business units, ensuring that data conforms to the enterprise’s established quality and consistency standards. Key responsibilities include defining and standardizing data elements, monitoring data quality, and collaborating with IT to resolve any technical challenges. Other key data roles are: Chief data officers (CDOs) lead the data strategy, ensuring data is treated as a valuable business asset. Their goal is to drive executive investment in data compliance, risk reduction, and value creation as data becomes a trusted driver of business outcomes.Data protection officers (DPOs) ensure organizational compliance with data privacy laws like GDPR and CCPA. They oversee the protection of personal data, such as that of customers or suppliers, processed during daily operations. DPOs must have direct access to senior leadership to fulfill regulatory requirements.Data architects design robust yet flexible data foundations that empower users to manage and enhance their own datasets. They ensure data is meaningful, business-driven, and aligned with organizational goals. Their priorities often reflect measurable business outcomes.Data engineers and developers design and maintain data pipelines, ensuring data quality and flow across complex systems. They aim to empower business users while managing access, security, and data product governance.Data scientists extract value from data pipelines to deliver actionable insights. They solve complex problems using statistics, mathematics, and computer science. Their expertise often includes data mining and predictive analytics.Business analysts identify trends, assess risks, and gauge business performance using BI tools like Tableau, Power BI, and Looker. They extract trusted insights from data pipelines and present them through clear, actionable dashboards. Data Policies and Procedures It establishes and enforces policies for how data is collected, stored, shared, and used. As enterprise central data management, Prescribes permitted and prohibited practices at every stage of the data lifecycleEnsures compliance with internal standards and external regulationsAssigns accountability for data stewardship and risk mitigationAligns day-to-day data handling with strategic business objectives Data Security Establishing proper security protocols helps in reducing the risk of data breaches and threats. It also safeguards sensitive information. Implementation of Encryption, access controls, authentication, and intrusion detection systems helps in protecting data across the lifecycle. Top open-source data security tools that are widely used include: Metasploit, OSSEC, OpenVAS, Snort, KeePass, ClamAV. These tools can be integrated into various security strategies to protect against a wide range of cyber threats. Metadata Management It helps in keeping track of data definitions, relationships, and structures. It’s essentially data about data. Metadata functions as the contextual glue that transforms isolated data points into coherent, actionable assets. It captures essential attributes covering: Creation timestampAuthorship and ownershipSource provenanceRelationships to other data elements Metadata strategy should: Adopt a centralized metadata catalog (e.g., Apache Atlas, Collibra)Automate metadata harvesting and lineage trackingIntegrate metadata-driven data quality checks into your pipelinesEstablish governance policies for metadata stewardship and versioningMonitor metadata KPIs like catalog adoption rate and lineage coverage to drive continuous improvement Leading open-source metadata management tools are Apache Atlas, Amundsen, Metacat Data Catalog, Open Metadata, and Marquez. Master Data Management Master data management is a process for ensuring the accuracy, consistency, and completeness of critical data elements, such as customer data and product data, etc. master data is standardized, matched, merged, enriched, and validated according to governance rules. Some of the open-source key players in the MDM area are Talend Open Studio for MDM, AtroCore, and Pimcore. Data Storage It helps determine where and how data will be stored within the enterprise data repository. It covers both structured and unstructured data sources, which include databases, data warehouses, and data lakes. The factors that determine data storage are Performance, scalability, and data retrieval requirements. Some of the key open-source players in the data storage area are Hadoop, LakeFS, Cassandra, and Neo4j. These tools provide scalability, robustness, and performance in managing large data and analyzing large datasets in various applications. Data Privacy and Compliance It ensures adherence to regulations and ethical considerations. Privacy implements controls to prevent unauthorized access and provides control over individuals' personal data. Regulatory frameworks such as the European Union’s General Data Protection Regulation (GDPR) and California’s Consumer Privacy Act (CCPA) impose stringent requirements on how businesses collect, process, and safeguard personal data. Data Metrics Management Defining and implementing robust business metrics and key performance indicators (KPIs) to quantify the enterprise-wide impact of data governance is critical to its success. These measures should be clearly articulated, inherently quantifiable, tracked longitudinally, and applied each year consistently to ensure comparability, accountability, and continuous improvement. Some of the metrics monitoring activities are, Aligning KPIs to strategic goals (e.g., data-quality gains, reduced time-to-insight, compliance rates, cost savings)Leveraging real-time dashboards for ongoing visibilityConducting annual KPI reviews to recalibrate targets and processes as the organization evolves Modern Data Architecture for AI A modern data architecture provides capabilities necessary for analytics, machine learning, generative AI, and autonomous systems. It enables enterprises to manage data as a strategic asset while ensuring governance, security, and scalability. The architecture is a unified, governed, AI-ready data foundation that enables trusted insights, intelligent automation, and autonomous decision-making through reusable data products, continuous observability, and embedded governance controls. The architecture is organized into two structural categories. The first five layers form the primary pipeline, the path data travels, from the moment it is created in a source system to the moment it produces a business outcome. The remaining three layers are cross-cutting disciplines that are applied continuously, at every stage, from ingestion through consumption. A modern AI-ready data architecture provides the infrastructure necessary for analytics, machine learning, generative AI, and autonomous systems. It enables organizations to manage data as a strategic asset while ensuring governance, security, and scalability. Figure 2: Enterprise Data Architecture For AI Data Sources This layer represents the full surface area of enterprise data — every system, channel, partner relationship, and unstructured artifact that generates information the organization can use. This layer groups the ecosystem into four major categories: Operational systems: The systems of record that run the business day-to-day: ERP, CRM, domain platforms, billing, and HR, etc. These remain the backbone of structured, transactional data.Digital channels: Web, mobile, API, and customer portal through which customers and employees interact directly with the enterprise. These channels are not purely a source; they also receive personalized or real-time data back through APIs.Partner ecosystems: B2B integrations, data exchanges, and marketplaces that bring external, third-party data into the enterprise's view.Unstructured and knowledge: Documents, email, video, knowledge bases, and ontologies. This category has grown in strategic importance because it is precisely the content that large language models and retrieval-augmented generation (RAG) pipelines depend on. Ingestion, Integration, and Orchestration This helps to move data from source into the platform reliably, securely, and in the right cadence, like batch, streaming, or on-demand. This layer comprises four capability areas, Data pipelines and orchestration: Engines that sequence and monitor data movement, paired with pipeline observability so failures and delays are visible before they become business problems.API management: Gateways, throttling, versioning, and security policy enforcement for every API-based integration, ensuring that data movement through APIs is controlled rather than ad hoc.Streaming and events: Event hubs and pub/sub infrastructure (e.g., Kafka-style platforms) that support event-driven integration for use cases where near-real-time movement is required.Data virtualization: Query federation that lets consumers query across multiple heterogeneous stores without first physically consolidating the data, reducing duplication and latency for enterprise usage. Core Data Platform (Analytics + AI) This is the heart of the architecture that acts as an AI-ready layer. It provides a unified storage and serving layer. This is the place where data lives and is made available for both traditional analytics and AI workloads from a single, governed foundation. Lakehouse and warehouse: It combines the flexibility of a data lake with the performance and semantic structure of a warehouse, including reusable semantic models that give consistent business meaning to raw tables.Operational data stores (ODS): Supports near-real-time reporting for use cases that cannot wait for a batch cycle.Vector and knowledge layer: Vector databases and ontologies that power agentic AI and semantic search are foundational to GenAI.Feature and model stores: Reusable features, a model registry, and model artifact storage, enabling machine learning models to be built, versioned, and reused consistently rather than recreated per project.Content and document stores: A repository that supports GenAI applications operating directly over enterprise content (contracts, policies, knowledge articles). AI, Analytics, and Decision Intelligence In this layer, the governed data is converted into insight, prediction, and increasingly autonomous action. Descriptive and diagnostic: BI, dashboards, and self-service analyticsPredictive and prescriptive: Machine learning models, optimization, and simulation GenAI and agentic AI: Copilots, task-oriented agents, and RAG pipelines that generate content to take bounded actions on the enterprise's own dataDecision intelligence: Composite decision flows that blend rules engines, analytics, and AI models into a single decision path Data Management and Semantics Layer Makes data trustworthy, findable, and consistently defined. This is applied continuously across every stage of the pipeline rather than as a single processing step. Enterprise data catalog: Technical and business metadata plus a data marketplace, such that stakeholders and systems can discover what data exists and what it means.Business glossary: Shared definitions, metrics, and domain vocabularies that prevent the classic problem of different business units calculating "revenue" or "active customer" differently.MDM and reference data: Golden records for core entities such as provider or product, eliminating duplication and conflicting versions of the truth.Data quality and profiling: Rules, scoring, and remediation workflows that continuously monitor and improve data fitness for use.Lifecycle management: Retention, archival, tiering, and deletion policies that keep the data estate compliant and cost-efficient over time. Agentic AI Governance, Security, and AI TRiSM Protects data and models with policy, privacy, identity, and full traceability. Policy-as-Code: Codified policies that are enforced programmatically rather than documentedLeast-privilege tool scope: Agents should operate with scoped function definitions rather than open-ended enterprise API access. Tools exposed to agents must enforce fine-grained parameter constraints AI TRiSM (Trust, Risk, and Security Management): Model risk assessment, explainability, fairness testing, and ongoing monitoring, addressing the risks introduced by AI/ML modelsIdentity delegation and impersonation: Enterprise agents must pass user identity context (OAuth 2.0 Token Exchange/On-Behalf-Of flow) down to underlying APIs. The agent must never inherit broader database permissions than the initiating user.Privacy and protection: PII/PHI classification, masking, and tokenization to limit exposure of sensitive data.Access and identity: RBAC/ABAC, fine-grained entitlements, and a Zero Trust posture, ensuring access is granted on a least-privilege basisLineage and observability: End-to-end lineage across data, models, and promptsPrompt/Context provenance and non-determinism audit: Every dynamic branch decision made by an agent must log its inputs, system prompts, retrieved context chunks, and seed parameters. This ensures that non-deterministic outputs can be audited post-hoc for compliance, debugging, and root-cause analysis during hallucinations or incorrect tool dispatches.Lineage granularity for vector and RAG workflows: Lineage models must extend beyond tabular source-to-target paths to map vector embedding lineage, tracing an agent’s final action back through the vector search embeddings, semantic chunking boundaries, and original unstructured document versions. Platform Engineering and MLOps/DataOps Dedicated engineering discipline. DataOps: CI/CD for data pipelines, including automated testing and deployment, bringing software-engineering rigor to pipeline changes.MLOps: CI/CD for models, including drift detection and automated retraining, so model performance is managed as an ongoing operational concern rather than a one-time deployment event.Platform engineering: Self-service portals, templates, and guardrails that let data and AI teams provision what they need quickly while staying within approved patterns.Infrastructure layer: Serverless compute, storage tiering, and cost management, ensuring the platform scales economically as usage grows. Business Consumption and Experience In this layer, the value is realized. The components and agents in this layer call back into the AI/Analytics layer in real time to inform what gets built upstream. Line-of-business applications: Domain applications, operations tooling, and customer service platforms through which employees and customers experience the businessCopilots and agents: Embedded copilots and agents inside applications and communication channelsAutomation and orchestration: Business process management (BPM), robotic process automation (RPA), and event-driven automation that act on insight without requiring manual interventionKPIs and value realization: OKRs, business outcome tracking, and benefit tracking that close the loop, measuring whether the solution is delivering value Benefits of Data Governance Enterprises with mature governance capabilities experience higher AI model accuracy, increased data trust, better decision-making, faster innovation cycles, improved compliance posture, reduced data duplication, and greater business agility. It also helps in: Ensuring consistent, uniform data across the enterprise, empowering smarter, more comprehensive decision supportEstablishing data integrity, data accuracy, completeness, trustworthiness, and dependability to achieve higher quality business decisionsHelping teams gain comprehensive decision support by enabling strong governance across the enterpriseDefining clear protocols for evolving data workflows; data governance helps in establishing agility and scalability for both the business and ITReducing duplication of effort and improving productivityMaking better-informed decisions with accurate and reliable dataIncreasing efficiency by introducing the ability to reuse data and data processesLowering the expenses of data management by implementing centralized control mechanisms and reducing the risk of data breachesReducing the volume of data collected and retained, optimizing data storage, and improving data management practicesEnhancing trust in the accuracy of data and the documentation of data-related proceduresEnsuring adherence to data regulations and supporting compliance with the EU’s GDPR, California Consumer Privacy Act (CCPA), Health Insurance Portability and Accountability Act (HIPAA), and the Payment Card Industry Data Security Standard (PCI-DSS) Conclusion Data governance is not a one-time activity, but it’s a journey. It is not optional but mandatory. It enables insight generation and informed decision-making. Effective data governance is a collection of processes, people, policies, standards, and metrics that ensure the efficient and effective use of data, enabling an enterprise to achieve its goals. It helps streamline operations, minimize data risks, enhance decision-making, drive innovation, create data policies, maximize data usage, and improve business efficiency and competitiveness. The modern AI data architecture is a unified, governed, and AI-ready foundation that turns enterprise data into trusted decisions and measurable business outcomes, with governance and AI risk management built in from the first byte rather than added at the end. By implementing data governance best practices, enterprises can ensure that they are managing their data to maximize its value. Acknowledgements The authors would like to thank Tricon Solutions LLC and Gspann Technologies, Inc for giving the required time and support in many ways in bringing up this article. Disclaimer The views expressed in this article/presentation are those of the authors, and Tricon Solutions LLC and Gspann Technologies, Inc. do not subscribe to the substance, veracity, or truthfulness of the said opinion.
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.
Running Apache Flink on a mainframe sounds odd at first. A modern stream processing engine on a platform most people call legacy? But take a closer look. It is not only possible. It might be a smart move for some of the largest financial institutions in the world. This post explores why some enterprises want Apache Flink on the mainframe, how it could work, and whether it is a brilliant innovation or a technical detour. Disclaimer: The views and opinions expressed in this blog are strictly my own and do not necessarily reflect the official policy or position of my employer. Mainframes Will Still Matter in 203X! A few months ago, I wrote about integrating Apache Kafka with mainframe systems. The blog covered various real-world examples across industries. The key message: Mainframes are still in use. In many organizations, they are not going away. They remain a central part of IT strategy, especially in banking, insurance, and the public sector. But they are not just legacy systems. Modern mainframes such as the IBM z17 offer the latest Telum II processor and support up to 64 terabytes of system memory. The z17 enables very large in‑memory workloads and faster processing for analytics and real‑time use cases. These systems also integrate on‑chip AI acceleration and optional AI‑focused hardware to support machine learning and real‑time decisions directly where mission‑critical data resides, while running modern Linux environments and container platforms. Some companies are still on the mainframe because they cannot easily migrate. But many others do not want to move away. Instead, they modernize around the mainframe. Apache Kafka and Flink play a key role in this journey. They enable a real-time data foundation that connects core systems with modern applications across environments. In future hybrid cloud strategies, this becomes even more critical. Kafka acts as the central nervous system, delivering the right data and context at the right time between on-prem mainframes and cloud-based AI services, including agentic AI and large language models. An event-driven architecture with hybrid streaming replication ensures business-critical decisions are made on fresh, reliable, and contextual information. Mainframe Migration Has Not Happened Ask any architect or CTO in banking. Mainframe migration has been on the roadmap for over two decades. Full replacement of core systems is still rare. However, it is important to distinguish between migration and offloading. Mainframe migration means shutting down mainframe workloads entirely and moving all applications and data to a new platform. There are many reasons: Risk is too highOrganizational resistance is strongMainframe skills are still needed but hard to findSystems are complex and deeply integratedThese applications run reliably and perform well Mainframe offloading, on the other hand, is much more common. It means moving selected workloads, queries, or processing tasks off the mainframe to more flexible and scalable platforms. This reduces load and cost on the mainframe while enabling innovation elsewhere. I have shared several real-world examples of offloading in action, using Kafka, IBM MQ, and Change Data Capture (CDC) tools like IBM IIDR or Precisely to synchronize and replicate data between mainframe systems and the cloud or distributed infrastructure in real time: Mainframe Offloading and Integration Examples. Because of this, many firms choose mainframe integration and a slow lift-and-shift leveraging the Strangler Fig design pattern over migration. Kafka is already helping. Flink is the next step. Apache Flink Meets the Mainframe: Unlikely Combo, Real Potential At first glance, Apache Flink and the mainframe seem like technologies from two different worlds. But combining them can unlock surprising value. What Is Apache Flink? Apache Flink is the leading open-source stream processing engine. It is designed to process high volumes of data continuously and in real time, rather than in batches. Flink is widely used to support use cases like fraud detection, customer personalization, operational monitoring, and data transformation at scale. Many of the largest tech companies and digital natives rely on Flink to process billions of events per day with low latency and high throughput. It supports both event streaming and batch workloads, but its true strength lies in real-time use cases. Here is an example of continuous stream processing leveraging Apache Flink together with OpenAI for Generative AI in real-time: Flink is built for modern environments. It runs natively on Kubernetes, integrates with Apache Kafka for real-time data ingestion, and is commonly deployed in public cloud, private cloud, or hybrid architectures. This makes it an ideal fit for enterprises looking to build fast, intelligent applications on fresh and contextual data. How to Run Apache Flink on the Mainframe? Yes, Apache Flink can run on the mainframe. In fact, it already does. I have already seen this deployed in a real-world environment. A large global financial institution is preparing to invest massively to expand its use of Apache Flink. Running Flink on IBM LinuxONE is a central part of that strategy. This is NOT a lab experiment, but a production-focused initiative. This bank already uses Kafka and Flink in production. Now they want to move Flink compute workloads onto the mainframe. The reason is simple. They already have unused compute on LinuxONE. Running Flink there is cheaper and easier to scale (for some companies) than scaling out other systems. The architecture is modern. IBM LinuxONE runs OpenShift. IBM LinuxONE is a high-performance, enterprise-grade server built on IBM Z architecture. It is designed to run Linux workloads with extreme reliability, scalability, and security. Unlike traditional mainframes focused on COBOL and legacy apps, LinuxONE is optimized for modern Linux applications. Flink is deployed in containers inside OpenShift's Kubernetes infrastructure, just like in any other cloud or data center. From a technical perspective, you need to build Docker images for the s390x architecture to run Apache Flink on IBM LinuxONE. In addition, components like RocksDB, which is used as a state backend in Flink, must be compiled for s390x to ensure full functionality. Why Put Apache Flink on the IBM Mainframe? This is a very valid question! Nobody would buy a mainframe just to run Flink on it. However, this approach offers several benefits for organizations that already own and operate mainframe infrastructure: Available compute resources on the mainframe.Consume data directly from mainframe sources (such as IBM MQ or other integration interfaces) and process it directly on the mainframe; or consume data from external sources such as a Kafka cluster running on x86 infrastructure, enabling flexible integration across hybrid environments.Lower total cost of ownership (TCO) regarding hardware and license costs compared to adding new external x86 servers and bi-directional integration pipelines.Simplified operations within a single, familiar environment.Benefit from IBM actively promoting LinuxONE and driving more workloads onto the platform. Mainframes are not only still alive. They are growing. IBM’s infrastructure business, which includes the mainframe, is doing very well. In Q3 2025, IBM reported 3.6 billion dollars in revenue for the infrastructure segment. That is 17 percent growth. IBM Z alone grew 61 percent. In Q2 2025, infrastructure revenue was 4.14 billion dollars, beating expectations by a wide margin. This is not legacy tech in decline. It is a platform in transformation. A New Chapter for Stream Processing and Mainframes Apache Flink running on the mainframe may sound unusual at first, but it reflects a broader shift in how enterprises think about modernization. The mainframe is not just a legacy system to replace. IBM Mainframe can fit hybrid cloud strategies, especially in highly regulated industries like banking and insurance. Apache Flink brings real-time intelligence. The mainframe brings performance, reliability, and unmatched security. Together, they offer a powerful combination for building fast, contextual, and mission-critical applications, without abandoning existing infrastructure. With Kafka as the backbone and Flink as the engine for real-time processing, organizations can connect mainframe systems with cloud innovation, including advanced AI workloads. This is not just about preserving the past. It is about extending and reusing trusted systems to meet the demands of the future. Enterprises that embrace this model can reduce risk, increase agility, and unlock new value from the heart of their operations.
Dashboards are everywhere. Business and IT teams use them to track metrics, visualize trends, and make decisions. But when working with real-time data from Apache Kafka, it’s not obvious how to connect dashboards to the stream or whether you should at all. The conversation often jumps to technical options like Flink SQL, Kafka Streams Interactive Queries, or Confluent's TableFlow. Others try to build interactive dashboards directly on top of Kafka topics using a JDBC connector into a database and a Business Intelligence tool. But that only makes sense once the actual goal is clear. What is the business trying to do with the data? Dashboards are not always the right tool. Automation, smart agents, or process intelligence often deliver more value. Let’s unpack the bigger picture. This blog post breaks down the different types of queries on Apache Kafka data, when dashboards make sense, and why a context engine often plays a key role. Why Dashboards — And When Not To Use Them Dashboards give people visual access to data. They support decisions, reporting, and oversight. But not all data needs to be visualized. Dashboards make sense when: Business users want a regular view of changing dataTeams need to investigate operational metricsThere is a requirement for manual filtering and inspection But in many cases, dashboards are not the best answer. For example A machine overheating should trigger an alert, not wait for someone to look at a graphA fraud detection system should act instantly, not visualize the anomalyAn AI agent monitoring supply chains should get structured context, not a dashboard snapshot In these scenarios, dashboards are a fallback. The real need is action or automation, not visualization. This is where agentic AI and process intelligence come into play. AI agents require structured, fresh context. They do not use dashboards. They consume streaming data, apply logic or reasoning, and trigger downstream actions. Dashboards might still be used to audit what happened but not to drive the process itself. So before jumping into dashboard tools, first ask: Is this data for a human to observe or a system to act on? Foundations First: Apache Kafka, Event Streaming, Data Products, and Governance Apache Kafka is the core of modern event-driven architecture. It enables systems to stream events in real time, such as customer interactions, machine signals, backend transactions, or system logs. Unlike batch pipelines, event streaming allows continuous data flow across the business. This supports responsive applications, automation, and real-time analytics. But fast data is not enough. Real-time value depends on reliable data. That’s why many teams now treat Kafka topics as data products. Each stream should have a clear owner, a defined schema, and a contract between producers and consumers. Schemas must be versioned and validated. Metadata must be consistent and available. Lineage, access control, and quality checks are critical to avoid downstream errors. Without this foundation, queries will return incorrect results, and automation may act on bad signals. Governance, schema control, and product thinking are not extras. They are required to build trustworthy systems on streaming data. Three Kinds of Queries for Apache Kafka Events If a dashboard is needed, the next step is understanding the type of query behind it. This helps define the right technical setup. Operational Queries These are fully automated. They respond to events and trigger actions. Think of them as the nervous system of an application. They are built directly into stream processing applications using Apache Flink or Kafka Streams. The logic is reactive and runs continuously. These systems are part of mission-critical operations. They must be highly available, fault-tolerant, and operate with minimal latency. Any downtime or delay can disrupt core business processes. A modern data streaming platform that augments Kafka and Flink with on-the-fly table serving, snapshot queries, and a context engine helps close this gap between streaming and interactive exploration. Example use cases include raising alerts on thresholds, aggregating orders for reporting, or triggering workflows. These systems should not rely on dashboards. Explorative Queries These are used by people to explore the data. They are ad hoc, flexible, and interactive. This type of query is difficult to support directly on Apache Kafka. Kafka is optimized for high-throughput event streaming and acts as an immutable event log. It provides a durable persistence layer and decouples producers from consumers, which makes it ideal for data pipelines and ensuring data consistency across real-time and batch systems. However, it is not designed for indexed lookups or ad hoc filtering across large datasets. Kafka does not offer queryable storage, secondary indexes, or snapshot consistency, all of which are essential for interactive exploration. Flink can process the data, but it does not offer indexed access. That makes joins or drilldowns inefficient without an external engine. Exploratory queries are often run in SQL workbenches, BI tools like Superset, or analytical engines like Druid and ClickHouse. They are useful for finding anomalies, trying out new logic, or investigating correlations. They require indexing, snapshot consistency, and historical access. Example use cases include joining marketing and sales events to find conversion patterns, analyzing user journeys through digital platforms, or testing new business rules across historical data. These queries typically require interactive tools and should not rely on stream processing systems alone. Monitoring Dashboards This use case is simpler but more common. The goal is to display filtered, consistent, and up-to-date data to end users. It does not involve complex joins or deep exploration. Instead, dashboards show metrics from recent data, business KPIs, or precomputed aggregations. Tools used here include Power BI, Grafana, or custom frontends connected to Flink or TableFlow. Dashboards in this case should be thin and rely on upstream systems for logic. Example use cases include showing live production status on a factory screen, displaying transaction volumes in a finance dashboard, or visualizing the health of streaming pipelines for operations teams. These dashboards are read-only and should not contain business logic. What Businesses Really Need Today While use cases vary, a few patterns repeat across industries. These needs can guide architecture decisions. Lightweight dashboards with filtering but no complex joins: Power BI and Grafana are the most common tools. Used for message tracing, monitoring, and status overviews. Users prefer querying externally instead of importing data.Real-time data that stays up to date: Dashboards refresh automatically. Data is pushed from Flink or precomputed topics. Materialized views support this, but changing schema can cause frontend problems.Business logic belongs upstream: Dashboards should not do computation. Flink or Kafka Streams handle the logic and prepare the data.Integration with ML models and agents: Dashboards may show results from predictions or scoring models. These are often trained ML models, not LLMs. Model drift monitoring is gaining interest. LLMs are still early stage in these setups.Protocol-agnostic connectors: REST, WebSocket, MQTT, JDBC — all needed. Most organizations expect flexible integration. Sink connectors alone are often not enough. APIs with query parameters are common requests. The Context Engine: Serving Dashboards and AI Agents from Apache Kafka Events A powerful pattern is the context engine. It connects Kafka streams to dashboards and AI systems by offering real-time, structured, and indexed access to data. It works like this Flink or Kafka Streams process raw Kafka topicsOutput data flows to context topicsA service builds indexed views of relevant business objectsDashboards and agents query those views through an API This setup creates a reliable source of truth. Business logic stays in the stream. The context engine focuses on enrichment, access control, and exposing views. For AI agents, this API layer usually follows the Model Context Protocol (MCP), which is becoming the de facto interface for connecting agents to structured enterprise data. Dashboards, in contrast, are typically served from materialized views in cache or in-memory databases, or directly through REST APIs optimized for low-latency reads. Agentic AI systems benefit directly. They consume these views as context to make decisions in real time. Instead of querying raw data or relying on stale batches, they get structured signals. Generative AI also benefits, using the same views as grounding data. Dashboards and AI agents both rely on fresh, accurate context. A context engine provides that bridge. Start With the Use Case, Not the Tool The right dashboard architecture does not start with a tool choice. It starts with business needs. Ask the right questions: What decisions or actions should this data support?Is the goal observation or automation?Does the user need filtering, drilldowns, or live KPIs?How fresh must the data be?Can the logic run upstream, or must it remain flexible? These answers will guide the setup. Sometimes a simple Power BI dashboard is enough. Other times a context engine or Flink job is required. In many cases, a dashboard is just the user interface to something much more powerful running behind the scenes. Of course, even when the focus is on business outcomes, a tool still has to be selected. That decision should follow the use case, not drive it. There are many options. Some teams prefer code-driven frameworks that give full control and allow deep integration with APIs and AI agent interfaces. Others choose no-code or low-code tools with prebuilt widgets so business users can create interactive views quickly. Each option comes with trade-offs in flexibility, governance, scalability, and integration. Exploring these tooling choices in depth would fill an entire chapter on its own. The key message here is simple: start with the outcome. The tool is an implementation detail. Build for the decision, not for the visualization. That is how streaming data creates real business value.
For the past two decades, most enterprise data engineering systems have been built on one default assumption: People understand the system. The system executes the pipeline. Engineers understand the business context, break a requirement into steps, write SQL, Spark jobs, shell scripts, synchronization tasks, and scheduling workflows, and then let the system run them. The scheduler does not need to understand the business. The sync engine does not need to understand the metric. It only needs to execute the predefined flow reliably. That model supported the era of data warehouses, data lakes, BI reporting, and batch scheduling very well. But now that assumption is starting to break down. Enterprise data systems are becoming more complex in every direction: More data sources.Longer pipelines.Stronger real-time requirements.Faster business changes.More conflicting metric definitions.More AI application data, model feedback data, vector indexes, and unstructured content. In this environment, enterprises do not just need more pipelines, and they do not just need a better Copilot that can write SQL faster. They increasingly need a Data Engineering Agent that can understand the system, plan tasks, call tools, validate outcomes, and accumulate experience over time. In that shift, Apache SeaTunnel becomes especially important. Because in the agent era, it is not enough for a system to "think." It also has to connect to real data sources, capture changes, execute synchronization, process incremental updates, preserve consistency, and move data to target systems in a reliable and cost-effective way. In other words: The agent understands the goal and plans the action. SeaTunnel turns that action into real, reliable, and recoverable data movement. That is why SeaTunnel is well positioned to become a core execution foundation in the evolution from ETL, ELT, and EtLT to agent-driven data engineering. ETL to ELT: The First Major Shift Traditional ETL is straightforward: Extract data from the source.Transform it in an intermediate layer.Load the processed result into the target system. This model fit the early data warehouse era well. At that time, data sources were relatively limited, the pipeline was easier to understand, and compute resources were more centralized. Enterprises wanted to clean the data, standardize the structure, and define the core logic before loading data into the warehouse. At its core, ETL is a deterministic pipeline model. Its key assumption is: People define the process in advance. The system executes the process. Later, with the rise of cloud warehouses, data lakes, lakehouse architectures, and elastic compute, ELT became more popular. ELT changed the order: ExtractLoadTransform inside the target platform. Instead of transforming everything before loading, enterprises started moving raw or near-raw data into a unified storage layer first, then using the target platform's compute power for downstream modeling and analytics. ELT solved several ETL limitations: It reduced upfront processing complexity.It preserved more original data.It gave analysts and modeling teams more flexibility later. But ELT also created a new problem. If all transformation is delayed until after loading, then dirty source data, schema drift, type mismatches, CDC events, privacy fields, and format inconsistencies all arrive directly in the target system. That might be acceptable in simple batch scenarios. It becomes much more expensive in real-time synchronization, CDC, multi-table sync, lakehouse ingestion, SaaS API ingestion, and AI-oriented data engineering. That is where a third pattern becomes more useful: EtLT. Why EtLT Matters EtLT is not just a compromise between ETL and ELT. A more useful way to understand it is: Extract -> lightweight transform -> Load -> semantic Transform That means: Extract the data.Apply the minimum engineering transformations required to make the data usable.Load it into a unified data foundation.Apply business-level and semantic transformation later. The key idea is the distinction between lowercase t and uppercase T. Lowercase t is not heavy business modeling. It is the engineering work that must happen before data enters the platform safely and consistently, such as: Field projectionType mappingFormat normalizationPrimary key or partition field handlingSensitive field maskingCDC event conversionMulti-table routingSchema evolution handlingPre-ingestion quality validationOne-read, multi-write patternsRate limiting and parallelism control. These transformations should not always be postponed to the target system. Otherwise, the lakehouse or warehouse becomes full of inconsistent, weakly governed, and semantically unclear raw data. At the same time, lowercase t should not try to absorb all business logic. Complex business definitions, KPI semantics, subject-area modeling, and cross-domain aggregation still belong to uppercase T, which should happen in the warehouse, lakehouse, semantic layer, or metric layer. That is the value of EtLT: Standardize the data engineering layer before loading, then apply business semantics after loading. This is exactly the place where SeaTunnel fits naturally. Its Source, Transform, and Sink architecture is well suited for the lowercase t in EtLT. It can connect heterogeneous systems, apply lightweight transformation during movement, handle CDC, adapt schemas, route multiple tables, and write the result into the target platform. In an EtLT architecture, SeaTunnel is not just a data mover. It becomes the data integration runtime that prepares data before it enters the unified data foundation. Why Traditional ETL Starts to Struggle Traditional ETL is built for relatively stable pipelines. You write the rules, draw the DAG, schedule the tasks, and fix failures when they happen. But modern enterprise data environments are no longer that simple. Today a single enterprise may operate across: OLTP databasesKafka streamsCDC pipelinesSaaS APIsObject storageLogs and eventsLakehouse platformsReal-time OLAP systemsVector databasesAI interaction logsModel output datasets. The problem is not only that there is more data. The data is also more fragmented, more heterogeneous, and more real-time. Pipeline length is another issue. A single business metric may depend on dozens of tables, multiple layers of wide tables, several business domains, and a long chain of definition changes. At that point, many enterprises no longer struggle with "Can we build the workflow?" They struggle with "Can anyone still explain the whole pipeline end to end?" This is where traditional ETL shows a structural limitation. One renamed field can break hundreds of tasks.One changed enum can silently shift multiple core metrics.One incorrect incremental logic branch can pollute an entire downstream analysis chain. The scheduler can tell you that a task failed. It usually cannot tell you why that task matters. The sync tool can move the data. It usually cannot tell you which business metric is now at risk. The engineer can fix the script. But only if that engineer can first rebuild the missing context. So the real weakness of traditional ETL is not just performance or reliability. It is that: It can execute the process, but it does not understand the system. Why Copilot is Not Enough Many teams first bring AI into data engineering through Copilot-style workflows: Generate SQLComplete Spark codeDraft YAMLProduce test samples. These capabilities are useful. They improve local productivity. But they do not solve the deepest problem in enterprise data engineering. Because the hardest part of data engineering is rarely just code generation. It is system understanding. Copilot can help generate a SQL statement, but it does not know the real business meaning of the field. It can help draft a synchronization task, but it does not know which downstream metrics will be affected by a schema change. It can help generate a scheduler config, but it does not know whether the change breaks historical consistency or recovery semantics. What enterprises actually struggle with includes: Lineage reasoningDependency analysisSemantic understandingMetric governanceRisk estimationImpact analysisIncremental recovery. These are not just autocomplete problems. So enterprises do not only need an AI IDE. They increasingly need an agentic data engineering system that can understand the target, decompose tasks, call engineering tools, and verify the result. The Real Shift: From Pipeline to Agent If we keep only one conclusion, it is this: Traditional ETL is "people define the process, systems execute the process." Agentic data engineering is "people define the goal, systems generate the process." That is not a slogan. It is a change in how work is organized. In the traditional model, engineers design the task chain first, configure Source, Transform, and Sink, and then let the scheduler execute the pipeline. The system faces a fixed process. In the agent model, the input may only be a business goal. For example: Add a new gross margin metric for orders and keep it aligned with the finance definition. Traditionally, the engineer must: Identify relevant data sources.Read table schemas.Inspect lineage.Design transformation logic.Configure sync and scheduling jobs.Add quality checks.Run regression validation. In an agent-oriented workflow, the system should be able to generate a sequence of actions around the goal: Identify the affected business entities.Discover candidate data sources.Analyze upstream lineage.Decide whether the job belongs to ETL, ELT, or EtLT.Generate or update the SeaTunnel synchronization task.Configure full-load or CDC mode.Apply lightweight transformation.Write the result into the warehouse or lakehouse.Trigger data quality validation.Evaluate downstream impact.Present the result for human confirmation. That is the real difference. The breakthrough is not "AI wrote a SQL statement for me." The breakthrough is: The system starts generating engineering actions from a business goal. But this immediately raises a critical question: When the agent plans a data action, who executes it reliably? That is exactly where SeaTunnel becomes essential. Apache SeaTunnel in the Agent Era: The Data Integration Execution Layer An agent cannot stop at reasoning and recommendations. If a Data Engineering Agent decides that a table should be synchronized, a CDC job should be adjusted, a broken pipeline segment should be replayed, or a data slice should be reloaded into the target system, it needs a stable and observable execution layer to carry out that decision. That execution layer needs several core capabilities. 1. It Must Connect to Many Kinds of Data Sources Enterprise data systems are inherently heterogeneous. An agent cannot live in a world with only one database or one file system. It needs to connect to MySQL, Oracle, PostgreSQL, SQL Server, Kafka, Hive, Iceberg, Doris, ClickHouse, StarRocks, Elasticsearch, S3, HDFS, MongoDB, and many other systems. SeaTunnel's connector architecture is designed for exactly this kind of environment. It abstracts Source, Transform, and Sink through a consistent plugin model so heterogeneous systems can be integrated in a unified way. 2. It Must Support Batch, Streaming, CDC, and Large-Scale Synchronization The agent era does not run on a single data movement pattern. It needs: One-time full migrationContinuous CDCOffline batch movementReal-time synchronizationSingle-table syncMulti-table or database-level sync. SeaTunnel is valuable here because it is not just a script wrapper for ETL. It is a real data integration runtime that can support full load, incremental sync, real-time processing, CDC, and multi-table movement in the same ecosystem. 3. It Must Handle the Lowercase t in EtLT Agentic systems do not need every business transformation to happen inside the sync layer. But they do need the sync layer to complete the minimum engineering transformation required to make the data trustworthy and usable before it lands in the platform. SeaTunnel's Transform layer is a strong fit for: Field mappingType conversionFilteringColumn projectionData maskingRoutingSimple reshaping. That is exactly the role of the lowercase t in EtLT: Do not overload the movement layer with heavy business modeling, but make the data governable and ready for the next stage. 4. It Must Provide Consistency, Fault Tolerance, and Recovery An agent can decide that a broken link should be replayed. But replay only matters if the underlying system can recover correctly. The execution layer still needs checkpointing, failure recovery, state handling, restart behavior, and strong delivery guarantees where needed. A reasoning layer without a reliable execution layer becomes a planner without hands. That is why execution quality still matters as much as intelligence. What the Future Stack Starts to Look Like If we look one step ahead, enterprise data engineering increasingly resembles a layered operating system rather than a collection of disconnected pipelines. In that stack: The semantic layer defines the business model.Metadata provides structure and context.Memory accumulates operational experience.The planning layer turns goals into actions.The execution layer performs synchronization, CDC, movement, replay, and recovery. SeaTunnel belongs to this execution layer. That placement is important. The future is not "put a large language model on top of ETL." The future is a coordinated system where reasoning and execution are separated clearly: The agent decides what should happen.SeaTunnel ensures that it actually happens in a reliable way. The Evolution in One Sentence ETL built data pipelines. ELT moved raw data into a unified platform first. EtLT rebalanced pre-load engineering standardization and post-load semantic modeling. The agent era pushes data engineering one step further: From fixed pipelines to goal-driven systems. In that world, SeaTunnel is not just a synchronization tool. It becomes a practical execution foundation for agentic data engineering. Agents make data systems understand goals. EtLT makes ingestion more controllable. SeaTunnel turns those goals into reliable data engineering actions. That is the deeper change now happening across enterprise data engineering.
A large API response becomes a client problem long before it becomes a network problem. A browser can receive hundreds of megabytes and still become unresponsive while buffering bytes, parsing one enormous JSON document, retaining duplicate object graphs, and rendering too much state on the main thread. The reliable solution is not a larger timeout. It is to stop treating the response as a synchronous document and start treating it as a durable, observable job whose data arrives in bounded pieces. Browser streams support incremental consumption and backpressure, while background workers allow long-running processing to remain independent of user-interface scripts. The Response Becomes a Job, Not a Payload The public API should acknowledge work quickly and return a stable job identifier rather than hold an HTTP connection open until every upstream page has been fetched. A 202 Accepted response establishes that contract without implying completion. The client can then subscribe to progress events, request a partial view, or retrieve a final artifact when the job reaches a terminal state. RFC 9110 defines 202 Accepted specifically for requests accepted for processing when processing has not necessarily completed. Java @PostMapping("/reports") public ResponseEntity<JobAccepted> create(@RequestBody ReportRequest request) { String jobId = UUID.randomUUID().toString(); workflowClient.start(reportWorkflow::run, jobId, request); return ResponseEntity.accepted() .header("Location", "/reports/" + jobId) .body(new JobAccepted(jobId, "QUEUED")); } This endpoint performs no large download or expensive transformation. It creates an addressable unit of work and returns immediately. The browser remains responsive because the initial response is tiny, while server capacity is protected from long-lived request threads. The job record should expose states such as queued, fetching, indexing, ready, failed, and canceled, with progress kept monotonic and coarse enough to remain trustworthy. Temporal Owns the Long-Running Control Flow Temporal fits the control plane because Workflow state survives process crashes and worker restarts, while failure-prone operations such as remote API calls belong in Activities with explicit timeouts and retry policies. Temporal documentation distinguishes deterministic Workflow logic from non-deterministic Activities and provides retry, timeout, heartbeat, and message-passing mechanisms for long-running execution. Java @WorkflowMethod public ResultRef run(String jobId, ReportRequest request) { String cursor = null; int sequence = 0; do { PageRef page = activities.fetchAndStore(jobId, cursor, sequence); activities.publishChunkReady(jobId, page); cursor = page.nextCursor(); sequence++; } while (cursor != null && !canceled); activities.buildIndex(jobId); activities.publishCompleted(jobId, sequence); return new ResultRef(jobId, sequence); } @SignalMethod public void cancel() { canceled = true; } Only references and counters should cross Workflow boundaries. Passing raw pages through Temporal causes every Activity input and result to accumulate in Event History. Temporal warns that large histories increase Workflow Task latency, documents a 50 MB or 51,200-event history limit, and recommends external storage plus Continue-As-New for large or long-running executions. The response body therefore belongs in object storage, while Temporal retains keys, checksums, cursors, and status. The fetching Activity should checkpoint often enough to support retries without restarting the transfer. Heartbeat details can carry the last committed cursor or byte range. Temporal recommends heartbeats for long-running Activities because missed heartbeats can trigger failure detection and retry. Java public PageRef fetchAndStore(String jobId, String cursor, int sequence) { UpstreamPage page = upstream.fetch(cursor); String key = storage.put(jobId + "/" + sequence, page.bytes()); Activity.getExecutionContext().heartbeat( new FetchCheckpoint(sequence, page.nextCursor()) ); return new PageRef( key, sequence, page.nextCursor(), page.sha256() ); } Kafka Carries Facts, Not Giant Documents Kafka is most effective as the event backbone, not as a substitute for object storage. Events should describe what happened and point to durable data, ChunkStored, ChunkIndexed, JobProgressed, JobCompleted, or JobFailed. Kafka enforces record-size limits at both producer and broker levels, so pushing multi-megabyte fragments into records creates brittle configuration coupling and expensive retries. Every event should use jobId as the key. Kafka partitions are ordered logs, and records sharing a key normally land in the same partition, preserving per-job sequence while allowing unrelated jobs to scale across partitions. Consumer groups distribute partitions across workers and rebalance them when membership changes. Java public void publishChunkReady(String jobId, PageRef page) { ChunkReady event = new ChunkReady( jobId, page.sequence(), page.storageKey(), page.sha256() ); kafkaTemplate.send("report-events", jobId, event); } Duplicate delivery must be assumed at every boundary. Kafka producer idempotence prevents duplicate writes caused by producer retries when compatible acknowledgment and in-flight settings are used, but downstream side effects still require idempotent consumers. An indexer can enforce uniqueness with (jobId, sequence, checksum) and commit its database transaction before acknowledging the Kafka offset. Backpressure should be expressed through bounded concurrency rather than hidden in memory. An Activity can publish one stored chunk at a time, while indexer lag indicates downstream pressure. Temporal can pause between pages when lag crosses a threshold, or consumers can scale until partition count becomes the limit. The Client Receives Progress and Bounded Content Server-sent events are sufficient when communication is primarily server-to-client. The protocol uses text/event-stream, keeps a persistent HTTP connection, and represents each notification as a small text block. A projection service can consume Kafka events, maintain the latest job state, and expose a resumable stream using application event IDs Java @GetMapping( value = "/reports/{jobId}/events", produces = MediaType.TEXT_EVENT_STREAM_VALUE ) public Flux<ServerSentEvent<JobEvent>> events( @PathVariable String jobId) { return eventProjection.stream(jobId) .map(event -> ServerSentEvent.<JobEvent>builder() .id(event.sequence().toString()) .event(event.type()) .data(event) .build()); } The client should render status changes and small previews, not append the full raw response into application state. When direct streaming is required, the Fetch API exposes the response body as a ReadableStream, allowing chunk-by-chunk processing rather than waiting for completion. Parsing should occur incrementally, with CPU-heavy decoding or transformation moved to a Web Worker, whose execution remains separate from user-interface scripts. Final delivery should usually be a paginated query API, a range-readable artifact, or a signed download URL. A giant JSON reconstruction endpoint merely recreates the original failure at the last step. RAG Turns Stored Volume Into a Useful Interface RAG becomes valuable after chunks are durably stored. Each chunk can be normalized, split along semantic boundaries, embedded, and indexed with metadata containing the job identifier, source sequence, object key, and byte range. The original RAG formulation combines parametric generation with retrieved non-parametric memory, grounding generation in selected passages rather than the entire corpus. Java @KafkaListener( topics = "report-events", groupId = "rag-indexers" ) public void onChunkReady(ChunkReady event) { if (index.exists( event.jobId(), event.sequence(), event.checksum())) { return; } byte[] payload = storage.get(event.storageKey()); chunker.split(payload).forEach(chunk -> index.upsert( event.jobId(), event.sequence(), chunk ) ); progress.markIndexed( event.jobId(), event.sequence() ); } The query path retrieves only the most relevant chunks and sends those bounded passages to the model. Raw object references remain attached so generated statements can link back to source material. RAG should not conceal incomplete ingestion; the query service must expose index coverage and reject complete-report requests until all expected chunks are indexed. Java public Answer answer(String jobId, String question) { List<Passage> context = index.search(jobId, question, 8); return generator.generate(question, context); } This layer changes the client experience from downloading everything before anything is useful to inspecting progress, searching partial results, and retrieving only relevant evidence. It also keeps model context bounded when the source response is extremely large. A Responsive System Is Built From Explicit Boundaries The essential boundary is simple: Temporal owns durable intent and recovery, Kafka distributes compact facts, object storage holds large bytes, RAG builds a searchable semantic view, and the client receives only bounded updates or explicitly requested slices. Each component solves a different failure mode, and none is forced to carry the complete response through an interface designed for small messages. The resulting architecture prevents UI freezes, survives retries and restarts, supports cancellation and replay, and makes large upstream results useful before a monolithic download could finish. Large-response handling becomes reliable when completion is modeled as a process rather than a payload.
I learned this lesson the hard way. We had a critical data pipeline running for over 3 hours every single day. The logic was perfectly clean. The overarching schema was explicitly right. There were absolutely no obvious memory leaks, and absolutely nothing looked fundamentally broken in the raw PySpark transformations. Then I finally checked the physical query plan. Under the hood, Apache Spark was quietly executing a massive Sort-Merge Join to merge a multi-terabyte fact table with a dimensional lookup table that was barely 50MB in total size. One single line of code changed—wrapping that exact tiny lookup table cleanly in a broadcast() hint—and the exact same analytic job plummeted from 3 hours to just 18 minutes. That was it. One word saved us hours of daily compute and massive underlying cloud FinOps costs. Wrong join types are financially devastating directly because they are completely silent. Spark will not throw an aggressive exception. Your pipeline will not explicitly fail. It will simply execute your logic confidently 10× slower than it architecturally ever needs to. Here is the exact mental model I exclusively use now every single time I write a distributed join in Apache Spark. TL;DR: Silent shuffle bottlenecks kill massive Spark performance. Always explicitly broadcast small tables (< 200MB), default natively to Sort-Merge for massive dual-sided joins, and actively, aggressively leverage AQE skew joins to fundamentally prevent heavy task skew. Check your physical plans! 1. The Small Table: Always Broadcast When actively joining a massive fact table heavily against a tiny dimension table (like cleanly mapping a primitive status_id logically to a status_name), globally shuffling the multi-terabyte fact table wildly across the distributed cluster is architectural suicide. The rule: If a table is reliably under 200MB, rigorously physically force a Broadcast Hash Join.The mechanism: Spark intelligently bypasses the massive network shuffle entirely. It simply naturally copies the tiny 50MB table directly into the RAM of every single native worker node, allowing them to map data logically and locally. The Implementation Python from pyspark.sql.functions import broadcast # Wrapping the small lookup table strictly natively in a broadcast hint enriched_df = massive_fact_df.join( broadcast(small_lookup_df), "customer_id", "left" ) 2. Both Sides Massive: Default to Sort-Merge If you are systematically actively joining two massive, multi-terabyte tables accurately together (e.g., dynamically merging historical transactions cleanly with historical web_sessions), you physically cannot organically broadcast data without instantly dynamically triggering brutal Out-Of-Memory (OOM) driver exceptions. The Rule: Default heavily unconditionally to the Sort-Merge Join.The Mechanism: This is Spark's absolute most robust, incredibly stable joining algorithm physically built for massive scale. Spark heavily and organically shuffles the massive data wildly across the cluster so that precisely matching keys uniquely land physically on the exact same nodes, fundamentally and strictly sort them, and actively, efficiently, and accurately merge them natively. It is technically slower than a pure broadcast, but it is incredibly beautifully resilient inherently at petabyte scale. The Implementation Python # No explicit hints structurally required. Spark will cleanly natively default seamlessly to Sort-Merge for massive large datasets. final_df = massive_transactions_df.join( massive_sessions_df, "user_id", "inner" ) 3. Highly Skewed Data: Enable AQE Skew Join In heavy enterprise datasets, physical data is rarely organically distributed perfectly evenly. Imagine an active e-commerce platform where the default "Guest Customer" cleanly accounts for physically 60% of all universal platform transactions. If you intelligently execute a naive Sort-Merge Join broadly on customer_id, one single isolated Spark executor will physically be forced systematically to exclusively process the entire massive 60% "Guest" chunk. The other 199 regular executors will efficiently and cleanly finish in seconds and sit completely idle while that one node globally grinds to a halt. The rule: Actively, safely leverage Adaptive Query Execution (AQE) dynamically to natively, beautifully split heavily skewed partitions dynamically.The mechanism: AQE actively, dynamically, and securely detects massively skewed partitions directly mid-flight, accurately splitting them cleanly into optimally smaller, incredibly uniform, reliable sub-partitions so they can be effectively and seamlessly processed rapidly and cleanly in parallel. The Implementation Python # Ensuring AQE and Skew Join optimization are physically aggressively cleanly enabled locally in the specific cluster config spark.conf.set("spark.sql.adaptive.enabled", "true") spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true") The Silent Hero: Adaptive Query Execution (AQE) The part most data engineers fundamentally miss: AQE has actually been turned ON by default since Apache Spark 3.2. This means Spark is reading actual statistical data mid-job. If you execute a Sort-Merge Join on a massive table that unexpectedly shrinks to 40MB after an aggressive .filter() clause, AQE will actively intercept the job mid-flight and organically auto-switch the execution directly into a lightning-fast Broadcast Join. You technically don't have to code anything for this to happen. It organically just happens. But you do have to verify it. You must explicitly verify that spark.sql.adaptive.enabled is active in your environment. You must actively understand exactly what it is doing—because when AQE occasionally guesses wrong (usually due to stale table statistics), you need to know precisely how to aggressively override it with manual hints. Conclusion Performance tuning in distributed compute engines fundamentally comes down to actively understanding the physical network shuffle. Check your explicit joins. Aggressively read your physical query plans (using .explain()). And never blindly trust default configurations at the enterprise level. What is the absolute worst join performance bottleneck you have ever hit in production? Let me know in the comments below!
Change data capture (CDC) pipelines look straightforward on paper: capture database changes, publish them to Kafka, and update downstream systems. The difficulty starts when events are duplicated, consumers restart, projections drift, or a team needs to replay months of history without corrupting the state it is trying to recover. A reliable CDC design has to account for those failure modes from the beginning. That means combining Kafka and Debezium with idempotent writes, deterministic projections, controlled replay workflows, reconciliation checks, and enough recovery evidence to explain what happened when something goes wrong. The architecture: The goal is not only to move inventory changes quickly. The goal is to make replay safe enough that operators can rebuild and explain the derived state after failure. This article builds one concrete pattern: The important detail is that replay safety is not a single feature. It is the result of several boring decisions lining up correctly. Data Model The data model should separate the aggregate state, the classification state, and the transaction history. PLSQL CREATE TABLE inventory_stock_on_hand ( sku VARCHAR(64) PRIMARY KEY, stock_on_hand BIGINT NOT NULL, updated_at TIMESTAMP NOT NULL ); CREATE TABLE inventory_bucket ( sku VARCHAR(64) NOT NULL, bucket_type VARCHAR(32) NOT NULL, location_id VARCHAR(64) NOT NULL, quantity BIGINT NOT NULL, updated_at TIMESTAMP NOT NULL, PRIMARY KEY (sku, bucket_type, location_id) ); CREATE TABLE inventory_transaction ( event_id VARCHAR(128) PRIMARY KEY, sku VARCHAR(64) NOT NULL, seller_id VARCHAR(64) NOT NULL, delta_quantity BIGINT NOT NULL, event_time TIMESTAMP NOT NULL, accepted_at TIMESTAMP NOT NULL ); CREATE INDEX idx_inventory_transaction_sku_time ON inventory_transaction (sku, event_time); CREATE INDEX idx_inventory_bucket_sku_bucket ON inventory_bucket (sku, bucket_type); The transaction table is the recovery anchor. If the availability projection drifts, the system needs a history to explain the projection. Do not rely only on the mutable aggregate table. inventory_stock_on_hand is useful for fast reads, but it is not enough for recovery. If the aggregate is wrong, it cannot explain how it became wrong. The accepted transaction history gives replay something durable to reason from. Ingestion Event Use an event ID that can survive retries and replay. JSON { "event_id": "mkt-evt-8f11a", "sku": "1231241", "quantity": 100, "operation": "I", "event_time": "2026-06-19T18:23:11Z", "seller_id": "seller-42" } The consumer should perform an idempotent write. One pattern is to insert the transaction first using event_id as the primary key. If the insert fails because the event already exists, skip the duplicate and emit a duplicate-suppression metric. Java public InventoryWriteResult apply(InventoryEvent event) { try { transactionRepository.insert(event.toTransactionRow()); } catch (DuplicateKeyException duplicate) { metrics.increment("inventory.duplicate_event"); return InventoryWriteResult.duplicate(event.eventId()); } stockRepository.incrementStockOnHand(event.sku(), event.quantity()); bucketRepository.incrementBucket(event.sku(), "SELLABLE", event.quantity()); return InventoryWriteResult.accepted(event.eventId()); } In production, the accepted transaction insert and the aggregate updates should be part of the same database transaction. A useful shape is: PLSQL BEGIN; WITH accepted AS ( INSERT INTO inventory_transaction ( event_id, sku, seller_id, delta_quantity, event_time, accepted_at ) VALUES ( :event_id, :sku, :seller_id, :delta_quantity, :event_time, now() ) ON CONFLICT (event_id) DO NOTHING RETURNING sku, delta_quantity ) INSERT INTO inventory_stock_on_hand (sku, stock_on_hand, updated_at) SELECT sku, delta_quantity, now() FROM accepted ON CONFLICT (sku) DO UPDATE SET stock_on_hand = inventory_stock_on_hand.stock_on_hand + EXCLUDED.stock_on_hand, updated_at = now(); COMMIT; That ON CONFLICT clause is not just a database convenience. It is part of the replay contract. It ensures that retrying the same business event does not apply the same inventory delta twice. Debezium Configuration Enable PostgreSQL logical decoding and configure Debezium to emit CDC topics for the inventory tables. JSON { "name": "postgres-inventory-connector", "config": { "connector.class": "io.debezium.connector.postgresql.PostgresConnector", "database.hostname": "<POSTGRES_HOSTNAME>", "database.port": "5432", "database.user": "<POSTGRES_USER>", "database.password": "<POSTGRES_PASSWORD>", "database.dbname": "<POSTGRES_DBNAME>", "topic.prefix": "inventory_source", "plugin.name": "pgoutput", "slot.name": "debezium_inventory_slot", "publication.autocreate.mode": "filtered", "table.include.list": "public.inventory_stock_on_hand,public.inventory_bucket,public.inventory_transaction", "snapshot.mode": "initial", "heartbeat.interval.ms": "10000", "tombstones.on.delete": "false", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable": "true", "value.converter.schemas.enable": "true" } } Debezium gives you history, but not recovery confidence. The confidence comes from how you key, project, replay, and reconcile that history. For replay work, track these connector facts in your runbook: Connector name and versionReplication slot namePublication name and included tablesSnapshot mode used for initial loadTopic prefixLast processed LSNConnector lagSchema history topic When a connector interruption happens, those details tell you whether you can resume normally, need a bounded replay, or need a new snapshot plus downstream reconciliation. Partition-Aware Routing The partition key should be chosen from the business ordering boundary. Java public class SkuPartitioner implements Partitioner { @Override public int partition( String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) { InventoryEvent event = (InventoryEvent) value; String orderingKey = event.getSku(); int partitionCount = cluster.partitionCountForTopic(topic); return Math.floorMod(orderingKey.hashCode(), partitionCount); } } Partitioning is not merely a throughput setting. If the projection depends on entity-local ordering, the entity belongs in the key. Kafka Streams Topology A simplified topology might rekey CDC records by SKU, materialize source tables, and compute availability. Java StreamsBuilder builder = new StreamsBuilder(); KTable<String, StockOnHand> stock = builder.table("inventory_source.public.inventory_stock_on_hand", Consumed.with(Serdes.String(), stockSerde)); KTable<String, InventoryBuckets> buckets = builder.table("inventory_source.public.inventory_bucket", Consumed.with(Serdes.String(), bucketSerde)); KTable<String, AvailabilityProjection> availability = stock.join( buckets, (stockRow, bucketRows) -> AvailabilityProjection.compute(stockRow, bucketRows), Materialized.<String, AvailabilityProjection, KeyValueStore<Bytes, byte[]>>as("availability-store") .withKeySerde(Serdes.String()) .withValueSerde(availabilitySerde) ); availability .toStream() .filter((sku, projection) -> projection.isPublishable()) .to("inventory.availability.v2", Produced.with(Serdes.String(), availabilitySerde)); The projection function should be deterministic. If replaying the same accepted history does not produce the same projection, the topology is not replay-safe. Recovery Contract Attach a Recovery Contract to the flow. YAML recovery_contract: flow: inventory-availability-projection tuple: "<H, O, I, F, S, Q, E>" history: source: - inventory_transaction - debezium.inventory_transaction order: key: sku idempotency: key: event_id duplicate_policy: skip_and_report function: name: compute_sellable_availability deterministic: true scope: supported: - by_sku - by_time_window - by_partition checks: - stock_on_hand_matches_transactions - sellable_quantity_non_negative - projection_event_time_valid evidence: - replay_scope - events_processed - duplicates_skipped - projections_changed - reconciliation_failures - confidence_status Treat this file as executable architecture documentation. A service should fail fast if the contract is incomplete for a critical flow. Java public final class RecoveryContractValidator { public void validate(RecoveryContract contract) { requireNonEmpty(contract.flow(), "flow"); requireNonEmpty(contract.history().source(), "history.source"); requireNonEmpty(contract.order().key(), "order.key"); requireNonEmpty(contract.idempotency().key(), "idempotency.key"); requireNonEmpty(contract.function().name(), "function.name"); requireTrue(contract.function().deterministic(), "projection must be deterministic"); requireNonEmpty(contract.scope().supported(), "scope.supported"); requireNonEmpty(contract.checks(), "checks"); requireNonEmpty(contract.evidence(), "evidence"); } private void requireNonEmpty(Object value, String field) { if (value == null || value.toString().isBlank()) { throw new IllegalArgumentException("Missing recovery contract field: " + field); } } private void requireTrue(boolean value, String message) { if (!value) { throw new IllegalArgumentException(message); } } } That validator does not make the system correct by itself. It prevents a more common failure: discovering during an incident that nobody defined the replay scope, idempotency key, or reconciliation checks. Replay Workflow Replay should be treated as a controlled workflow. Plain Text 1. Identify incident scope. 2. Select replay scope by SKU, time window, or partition. 3. Read authoritative history. 4. Rebuild deterministic projection. 5. Run reconciliation checks. 6. Emit recovery evidence. 7. Republish only if checks pass. The output should be an evidence report. JSON { "recovery_id": "rec-2026-06-19-001", "flow": "inventory-availability-projection", "events_processed": 1842, "duplicates_skipped": 17, "projection_rows_changed": 11, "reconciliation": { "stock_on_hand_matches_transactions": true, "sellable_quantity_non_negative": true, "projection_event_time_valid": true }, "confidence_status": "trusted" } A replay runner can keep the workflow explicit: Java public RecoveryEvidence replay(ReplayRequest request) { RecoveryContract contract = contracts.load(request.flow()); validator.validate(contract); ReplayScope scope = scopeResolver.resolve(request, contract); List<InventoryEvent> history = historyReader.read(contract.history(), scope); ReplayResult result = projector.rebuild(history, contract.function()); ReconciliationResult reconciliation = reconciliationRunner.run(contract.checks(), scope, result); RecoveryEvidence evidence = RecoveryEvidence.builder() .recoveryId(UUID.randomUUID().toString()) .flow(request.flow()) .scope(scope) .eventsProcessed(history.size()) .duplicatesSkipped(result.duplicatesSkipped()) .projectionsChanged(result.changedRows()) .reconciliation(reconciliation) .confidenceStatus(reconciliation.passed() ? "trusted" : "review_required") .build(); evidenceStore.write(evidence); if (request.publish() && reconciliation.passed()) { publisher.publish(result.projections()); } return evidence; } The replay runner should support dry runs. Dry runs let operators answer "What would change?" before republishing availability, billing, or detection outputs. Operational Metrics Track ordinary health and recovery confidence separately. Ordinary health: Consumer lagConnector lagTask restartsDLQ countEnd-to-end latency Recovery confidence: Replay durationReplay scope sizeDuplicate suppression countProjection rows changedReconciliation failuresConfidence status Example metric names: Plain Text inventory_ingest_events_total{result="accepted|duplicate|rejected"} inventory_cdc_connector_lag_seconds{connector="postgres-inventory-connector"} inventory_stream_projection_lag_seconds{topology="availability"} inventory_replay_duration_seconds{flow="inventory-availability-projection"} inventory_replay_events_processed_total{flow="inventory-availability-projection"} inventory_replay_duplicates_skipped_total{flow="inventory-availability-projection"} inventory_reconciliation_failures_total{check="stock_on_hand_matches_transactions"} inventory_recovery_confidence_status{status="trusted|review_required|failed"} Alert on disagreement, not only lag. A good pipeline can be caught up and still be wrong. YAML alerts: - name: InventoryProjectionReconciliationFailure expr: inventory_reconciliation_failures_total > 0 severity: page - name: InventoryReplayRequiresReview expr: inventory_recovery_confidence_status{status="review_required"} > 0 severity: ticket - name: InventoryConnectorLagHigh expr: inventory_cdc_connector_lag_seconds > 300 severity: ticket Reconciliation Queries Reconciliation should be executable, not just a diagram in a runbook. Start with invariants that are simple enough to automate. Example: Stock-on-hand should match accepted transaction deltas for a replay window. PLSQL WITH accepted_delta AS ( SELECT sku, SUM(delta_quantity) AS expected_delta FROM inventory_transaction WHERE accepted_at BETWEEN :from_time AND :to_time GROUP BY sku ), actual_delta AS ( SELECT sku, stock_on_hand - :baseline_stock_on_hand AS observed_delta FROM inventory_stock_on_hand WHERE sku = :sku ) SELECT a.sku, a.expected_delta, b.observed_delta, (a.expected_delta = b.observed_delta) AS matches FROM accepted_delta a JOIN actual_delta b ON a.sku = b.sku; Example: Sellable inventory should never be negative. PLSQL SELECT sku, location_id, quantity FROM inventory_bucket WHERE bucket_type = 'SELLABLE' AND quantity < 0; These queries are not academically exciting, but they are operationally powerful. They turn "the replay finished" into "the replay finished and the invariants passed." Replay Endpoint Sketch A replay workflow should be explicit and permissioned. One possible internal API: HTTP POST /internal/recovery/replay Content-Type: application/json { "flow": "inventory-availability-projection", "scope": { "type": "sku_and_time_window", "sku": "1231241", "from_event_time": "2026-06-19T18:00:00Z", "to_event_time": "2026-06-19T19:00:00Z" }, "dry_run": false, "requested_by": "sre-oncall", "reason": "projection drift after stream task restart" } The response should not just say 200 OK. JSON { "recovery_id": "rec-2026-06-19-001", "status": "trusted", "events_processed": 1842, "duplicates_skipped": 17, "projections_changed": 11, "reconciliation_failures": 0, "evidence_uri": "<RECOVERY_EVIDENCE_URI>" } The response is the operational artifact. It gives the team something to attach to an incident timeline and something to compare against later recovery runs. Tests for Replay Safety Replay safety should be tested before production incidents. Java @Test void replayingSameHistoryDoesNotChangeProjectionTwice() { List<InventoryEvent> history = List.of( event("evt-1", "SKU-1", 10), event("evt-2", "SKU-1", -2), event("evt-1", "SKU-1", 10) // duplicate ); AvailabilityProjection first = projector.replay(history); AvailabilityProjection second = projector.replay(history); assertThat(first).isEqualTo(second); assertThat(first.sellableQuantity()).isEqualTo(8); assertThat(first.duplicatesSkipped()).isEqualTo(1); } Also test late events, schema versions, partition rebalance, connector restart, and partial replay by entity. If replay is part of your recovery model, it deserves the same test discipline as the happy-path pipeline. Add failure injection tests that mirror production recovery: Java @Test void lateEventTriggersReviewWhenItChangesPublishedAvailability() { ReplayScope scope = ReplayScope.forSkuAndWindow( "SKU-1", Instant.parse("2026-06-19T18:00:00Z"), Instant.parse("2026-06-19T19:00:00Z") ); history.append(event("evt-1", "SKU-1", 10, "2026-06-19T18:01:00Z")); history.append(event("evt-2", "SKU-1", -3, "2026-06-19T18:59:00Z")); history.appendLate(event("evt-3", "SKU-1", -2, "2026-06-19T18:30:00Z")); RecoveryEvidence evidence = replayRunner.replay( ReplayRequest.dryRun("inventory-availability-projection", scope) ); assertThat(evidence.eventsProcessed()).isEqualTo(3); assertThat(evidence.projectionsChanged()).isGreaterThan(0); assertThat(evidence.confidenceStatus()).isEqualTo("review_required"); } Failure Injection Matrix Use a small matrix before every major release of the pipeline. Duplicate Event Injection: Send the same event_id twice.Expected evidence: duplicates_skipped > 0; no double-counted stock.Late Event Injection: Delay event arrival until after the projection has already published output.Expected evidence: late event count, changed projections, and review status if the output changes.Connector Pause Injection: Stop the Debezium connector for several minutes.Expected evidence: connector lag, replay scope, and reconciliation status.Offset Rewind Injection: Reprocess a known event range.Expected evidence: deterministic replay agreement.Schema Change Injection: Replay old and new schema versions.Expected evidence: schema versions recorded in the recovery evidence.Bad projection deploy Injection: Publish an incorrect derived state, then replay.Expected evidence: projections changed; reconciliation passes after rebuild. The point is not to create chaos for its own sake. The point is to practice the exact recovery motion before a real incident. Production Hardening Checklist Before relying on replay in production, confirm: The authoritative history has retention longer than the largest expected recovery window.The idempotency key is stable across producer retries.The Kafka partition key matches the business ordering boundary.The projection function is deterministic for the supported replay scope.The contract names every source topic, source table, check, and evidence field.The replay endpoint supports dry runs.Republish requires reconciliation success.Evidence is written to durable storage.Evidence records include schema versions and replay input bounds.Operators can find the runbook from the alert.The DLQ is treated as an input to recovery, not as the recovery plan itself. For high-value flows, make this checklist part of the architecture review. It is much cheaper to define replay semantics while designing the pipeline than to invent them under pressure. Common Mistakes Treating CDC topics as transient integration messages instead of durable recovery history.Choosing partition keys for infrastructure convenience rather than business ordering.Allowing stream processors to perform hidden non-idempotent side effects.Measuring lag but not correctness.Resetting offsets without a reconciliation plan.Assuming exactly-once semantics removes the need for recovery evidence. Conclusion Replay-safe CDC pipelines require more than Kafka, Debezium, and stream processing. They require explicit recovery semantics. Recovery Contracts give teams a compact way to define those semantics. Confidence-carrying replay gives operators evidence that the recovered state can be trusted. That is the difference between a pipeline that resumes and a platform that actually recovers.
In a previous article, we built a static supply chain graph in Neo4j using Apache Spark, with suppliers, warehouses, distribution centers, and retailers connected by shipping routes. That gave us a snapshot of the network at a point in time. In this article, we'll add the streaming layer: shipment events flow through Confluent Cloud Kafka in real time, land in Neo4j as enriched graph properties, and a live dashboard shows network health updating as events arrive. The full source code is available on GitHub. The Stack Each tool in the stack does what it does best: ToolRoleConfluent Cloud (free tier)Managed Kafka cluster and topicPython producer (Jupyter)Generates and publishes synthetic shipment eventsPython consumer (Jupyter)Consumes events and writes them into Neo4jNeo4j AuraDBGraph database storing the supply chain and shipment eventsPlotlyLive dashboard visualization One deliberate omission is that we aren't using the Neo4j Kafka Sink Connector, which is available as a managed connector on Confluent Cloud. That connector handles the consumer side automatically but carries a per-task hourly charge. For this article, we'll keep everything free by writing a Python consumer that does the same job. This also has a practical benefit: all the pipeline logic is visible in Python rather than hidden inside a managed connector configuration, which makes it easier to understand and adapt. The managed connector is a natural next step for production workloads. Setting Up Confluent Cloud Sign up at confluent.io and create a free cluster.Once the cluster is running, create a topic named shipment-events with 1 partition and default settings.Create an API key and secret under API Keys.Note the bootstrap server address from the cluster settings. Export these as environment variables in your shell: Shell export CONFLUENT_BOOTSTRAP_SERVERS=your_cluster.confluent.cloud:9092 export CONFLUENT_API_KEY=your_api_key export CONFLUENT_API_SECRET=your_api_secret Setting Up Neo4j AuraDB AuraDB is Neo4j's fully managed cloud database. A free tier is available with no credit card required. Sign up at console.neo4j.io/graphacademy.Create a new AuraDB Free instance.When the instance is created, download or note the credentials — the connection URI, username, and password. Neo4j only shows the password once, so save it somewhere safe.Once the instance is running, open the built-in Query tab and verify connectivity: MATCH (n) RETURN count(n). This should return 0. We are ready to load data. Before starting Jupyter, export the connection details as environment variables in your shell: Shell export NEO4J_URI=neo4j+s://xxxx.databases.neo4j.io export NEO4J_USERNAME=your_username_here export NEO4J_PASSWORD=your_password_here export NEO4J_DATABASE=your_database_name_here The Data Model Each shipment event represents a single status update for a shipment at a point in time. A shipment does not generate a sequence of events as it progresses — each event is an independent snapshot, which keeps the producer simple and the consumer stateless. The event structure is: JSON { "shipment_id": "c60eb761-f153-4840-8427-17fa9e34c56c", "supplier_id": "S013", "warehouse_id": "W005", "dist_center_id": "DC004", "retailer_id": "R025", "status": "delayed", "timestamp": "2026-08-04T12:57:15Z", "delay_minutes": 34 } Status follows one of four values — departed, in_transit, delayed or delivered, with a configurable delay probability. We use 15% delayed to make the dashboard interesting without overwhelming it. When the consumer writes an event into Neo4j, it creates a Shipment node and links it to the existing supply chain nodes via four relationship types: Cypher MERGE (sh:Shipment {shipment_id: $shipment_id}) SET sh.status = $status, sh.timestamp = $timestamp, sh.delay_minutes = $delay_minutes WITH sh MATCH (s:Supplier {id: $supplier_id}) MATCH (w:Warehouse {id: $warehouse_id}) MATCH (dc:DistributionCenter {id: $dist_center_id}) MATCH (r:Retailer {id: $retailer_id}) MERGE (s)-[:HAS_SHIPMENT]->(sh) MERGE (sh)-[:VIA_WAREHOUSE]->(w) MERGE (sh)-[:VIA_DIST_CENTER]->(dc) MERGE (sh)-[:DESTINED_FOR]->(r) MERGE on shipment_id means re-running the consumer never creates duplicate nodes. The Producer The producer notebook uses a fixed random seed to generate reproducible shipment events using IDs drawn from the existing supply chain and publishes them to Confluent Cloud via the confluent-kafka library: Python producer = Producer({ "bootstrap.servers": BOOTSTRAP_SERVERS, "security.protocol": "SASL_SSL", "sasl.mechanisms": "PLAIN", "sasl.username": API_KEY, "sasl.password": API_SECRET, "log_level": 0, }) Setting "log_level": 0 suppresses the librdkafka telemetry messages that appear otherwise. The producer supports both batch and continuous modes. For example: Python produce_events(num_events = -1) # stream continuously produce_events(num_events = 100) # publish exactly 100 events The display refreshes every PRINT_EVERY events using clear_output, showing the latest event and a running status breakdown — so the cell output stays manageable even when streaming thousands of events. The Consumer and Live Dashboard Rather than two separate notebooks, we combine the consumer and dashboard into a single pipeline. On each cycle, the loop: Polls Kafka for up to POLL_BATCH events and writes them to Neo4jQueries Neo4j for the current graph stateRebuilds and redraws the dashboardSleeps for REFRESH_INTERVAL seconds before repeating Rebuilding the full dashboard on every cycle is straightforward and works well at demo event rates. At higher throughput, a more efficient approach would be to update only the changed data rather than redrawing all eight panels on each refresh. The consumer uses its own Kafka group ID (supply-chain-dashboard) so it reads the topic independently, catching up on all existing events first before staying live: Python consumer = Consumer({ "bootstrap.servers": BOOTSTRAP_SERVERS, "security.protocol": "SASL_SSL", "sasl.mechanisms": "PLAIN", "sasl.username": API_KEY, "sasl.password": API_SECRET, "group.id": "supply-chain-dashboard", "auto.offset.reset": "earliest", "log_level": 0, }) The Live Dashboard The dashboard uses Plotly's make_subplots in a 4x2 grid, rebuilt on every refresh cycle using clear_output. Eight panels give a complete picture of network health: Row 1 – Overall Health Network status table: Total shipments, delayed count, delay rate, Kafka events consumed, refresh count, and any disabled nodesShipment status distribution: Donut chart showing the split between departed, in transit, delayed, and delivered, as shown in Figure 1 Figure 1. Shipment Status Distribution Row 2 – Warehouse View Delayed shipments by warehouse: Which warehouses are handling the most delayed shipments right nowWarehouse health score: A heatmap scoring each warehouse from 0.0 (everything delayed) to 1.0 (fully healthy), colored red through orange to green, as shown in Figure 2 Figure 2. Warehouse Health Score Row 3 – Origin and Destination Supplier performance: Which suppliers are generating the most delayed shipmentsRetailer impact: Which retailers are receiving the most delayed shipments — the downstream effect of any disruption Row 4 – Mid-Network and Flow Average delay by distribution center: Where in the middle layer delays are accumulatingShipment flow: A Sankey diagram (Figure 3) showing which suppliers are routing through which warehouses Figure 3. Shipment Flow - Suppliers to Warehouses The warehouse health score is the most immediately readable panel. The Cypher behind it computes the score directly in the graph: Cypher MATCH (sh:Shipment)-[:VIA_WAREHOUSE]->(w:Warehouse) WHERE w.active IS NULL OR w.active <> false WITH w.id AS warehouse, count(sh) AS total, count(CASE WHEN sh.status = 'delayed' THEN 1 END) AS delayed RETURN warehouse, round(1.0 - toFloat(delayed) / total, 3) AS health_score ORDER BY warehouse Simulating a Network Disruption One of the more compelling features of the graph model is how easy it is to simulate and visualize a disruption. Setting active = false on any node excludes it from the dashboard queries and the dashboard immediately reflects the simulated disruption on the next refresh cycle. We can do this before the dashboard starts: Python REMOVE_NODE = "W007" # mark this warehouse as inactive Or live, while the dashboard is running, using the Neo4j AuraDB Query tab: Cypher // Disable a node MATCH (n {id: "W007"}) SET n.active = false // Re-enable a node MATCH (n {id: "W007"}) REMOVE n.active // Check what is currently disabled MATCH (n) WHERE n.active = false RETURN labels(n)[0] AS label, n.id AS id Within 5 seconds, the dashboard reflects the change. The warehouse health heatmap shows the gap, the delayed shipments bar shifts to other warehouses as traffic reroutes, and the network status table shows the node as disabled. Re-enabling it and watching the metrics recover completes the disruption and recovery story. Standalone Operation At startup, the consumer notebook creates the supply chain nodes using MERGE. This operation is idempotent, so any existing nodes from the previous article are left unchanged. Note that this step creates nodes only — the relationships between supply chain nodes (supplier -> warehouse -> distribution center -> retailer) are assumed to exist from the previous article, or can be added separately if running this notebook in isolation. Python with driver.session(database = NEO4J_DATABASE) as session: for i in range(20): session.run("MERGE (:Supplier {id: $id})", id = f"S{i:03d}") for i in range(12): session.run("MERGE (:Warehouse {id: $id})", id = f"W{i:03d}") for i in range(10): session.run("MERGE (:DistributionCenter {id: $id})", id = f"DC{i:03d}") for i in range(30): session.run("MERGE (:Retailer {id: $id})", id = f"R{i:03d}") Gotchas and Lessons Learned Suppress librdkafka Logging Without "log_level": 0 in the producer and consumer config, Confluent's underlying librdkafka library prints telemetry messages to the cell output every time a connection is established. The messages are harmless. Suppress Neo4j Property Warnings Querying a property that does not yet exist on any node produces a GqlStatusObject warning from Neo4j for every query that references it. The active property falls into this category when no node has been disabled. The fix is one line to set notifications to "OFF" on the driver, as follows: Python driver = GraphDatabase.driver( NEO4J_URI, auth = (NEO4J_USERNAME, NEO4J_PASSWORD), notifications_min_severity = "OFF", ) Consumer Group Isolation Kafka distributes partitions across consumers in the same group, so each consumer processes only its assigned partitions. If we run multiple consumers using the same group ID against the same topic, each will only process a subset of the events. The dashboard uses supply-chain-dashboard as its group ID, and the tip is to run only one instance of this notebook at a time against the same topic and cluster. auto.offset.reset = earliest Without this setting, a consumer that starts after events have been published will miss everything that arrived before it connected. Setting earliest means the consumer always catches up on the full history of the topic before going live, which is essential if we stop and restart the dashboard mid-session. Clear Shipment Nodes Between Runs Each run of the consumer creates new Shipment nodes. Since the producer generates synthetic demo data, it's safe to clear these between runs; otherwise, successive runs would accumulate all historical shipments, and the dashboard counts would grow unbounded. The notebook clears all Shipment nodes at startup: Cypher MATCH (sh:Shipment) CALL (sh) { DETACH DELETE sh } IN TRANSACTIONS OF 10000 ROWS Summary We've built a real-time supply chain event streaming pipeline using Confluent Cloud Kafka and Neo4j. The producer generates synthetic shipment events continuously, the consumer writes them into the graph, and a live dashboard shows network health updating in near real-time. The disruption simulation — marking a node inactive mid-run and watching the dashboard respond — demonstrates one of the most compelling aspects of the graph model: the ability to ask structural questions about a network as it evolves. The same architecture adapts naturally to real logistics, IoT, or manufacturing event streams where understanding network structure matters as much as raw throughput. The full source code is available on GitHub.
Reliable orchestration for small language models depends less on model sophistication than on the durability of event flow and state. Under the assumptions used here — small model instances, little or no local state, Kafka as the event backbone, Temporal as the orchestration layer and durable state store, and Java as the runtime — the safest design is to treat model invocations as replayable side effects, Kafka as the transport and ordering substrate, and Temporal Workflow state as the canonical record of conversational progress. In that design, Kafka provides high-throughput append-only event delivery and partition-local ordering, while Temporal persists Workflow Event History and can replay execution after failures. Exactly-once semantics remain meaningful inside Kafka’s consume-transform-produce boundary when transactions and read_committed are used, but once processing crosses into external systems such as model APIs, durable activities, or databases, correctness comes from idempotency, deduplication, sequence checks, and reconciliation rather than from a global exactly-once guarantee. Assumptions The most productive baseline is a narrow one. Each conversation, task, or model session is keyed so related events land on the same Kafka partition, preserving order only where order actually exists: within one partition, not across the topic. Each workflow instance owns one conversational state machine, stores the minimal context needed to decide the next action, and invokes model calls through Temporal Activities so failures, retries, and timeouts are visible and durable. Large prompts, attachments, or long transcripts are not kept as incidental JVM memory because Temporal persists inputs and outputs in Event History and large histories degrade replay latency; those artifacts belong in external storage with durable references held in workflow state. Analysis The central engineering mistake in LLM orchestration is to confuse transport delivery with business completion. Kafka can guarantee at-least-once delivery by processing records before committing consumer offsets, and it can provide exactly-once behavior for Kafka-to-Kafka pipelines by atomically updating produced records and consumed offsets with transactions. Kafka’s own design documentation is explicit that the producer is the transactional component and that read_committed is advisable when aiming for exactly-once processing. The same documentation also makes clear why the guarantee weakens at system boundaries: once consumed data must be coordinated with an external state store or side effect, the problem becomes cross-system consistency rather than log delivery. In a Temporal-based model pipeline, that means Kafka should usually be treated as the durable ingress path, while Temporal owns the authoritative notion of whether an event was applied to a conversation state machine. That separation suggests a simple rule. Offsets are transport progress; workflow state is semantic progress. A consumer should therefore commit offsets only after handoff to a durable semantic owner. In this architecture, that owner is the Temporal workflow receiving a signal. Temporal workflows behave like stateful services that receive Signals, Queries, and Updates, and the platform persists Event History so a crashed worker can replay the workflow and resume from the last recorded event. Signal handlers are allowed to mutate workflow state, and blocking coordination can be expressed safely with Workflow.await. Activity retries are configured through ActivityOptions and RetryOptions, with heartbeat support for long-running calls. Java @WorkflowInterface interface ModelFlow { @WorkflowMethod void run(String sessionId); @SignalMethod void onEvent(ModelEvent event); @QueryMethod long lastAppliedSequence(); } private final ModelActivities activities = Workflow.newActivityStub( ModelActivities.class, ActivityOptions.newBuilder() .setStartToCloseTimeout(Duration.ofSeconds(20)) .setRetryOptions( RetryOptions.newBuilder() .setInitialInterval(Duration.ofMillis(250)) .setMaximumAttempts(5) .build()) .build()); private final NavigableMap<Long, ModelEvent> pending = new TreeMap<>(); private long nextSequence = 1; private ConversationState state = ConversationState.empty(); @Override public void onEvent(ModelEvent event) { pending.putIfAbsent(event.sequence(), event); } @Override public void run(String sessionId) { for (;;) { Workflow.await(() -> pending.containsKey(nextSequence) || state.closed()); if (state.closed()) break; var event = pending.remove(nextSequence); state = activities.applyEvent(state, event); nextSequence = event.sequence() + 1; } Workflow.await(Workflow::isEveryHandlerFinished); } @Override public long lastAppliedSequence() { return nextSequence - 1; } This workflow fragment does three important things at once. The @SignalMethod declares asynchronous event ingress, the @QueryMethod exposes durable progress for reconciliation, and the activity stub attaches retry policy directly to the state transition that may call a model endpoint or another dependency. The pending map is not a queue for throughput; it is a reordering guard. If Kafka redelivers a message or an upstream retry arrives out of sequence, putIfAbsent and the nextSequence gate prevent semantic duplication and preserve per-session causality. Finishing the run only after Workflow.isEveryHandlerFinished() avoids the Temporal-documented failure mode where a workflow completes or continues-as-new while a handler is still waiting on asynchronous work. The matching Kafka consumer must be deliberately conservative. Automatic commits are inappropriate because they advance transport progress in the background regardless of semantic application. Manual synchronous commits make the boundary explicit, and Kafka documents that committed offsets are the secure restart position, and that commitSync should write the next offset, meaning lastProcessedOffset + 1. The consumer is also not thread-safe, so per-partition in-order handling is easiest when one poll loop owns one consumer instance and performs durable handoff before commit. Java void pollLoop() { consumer.subscribe(List.of("model-events")); while (running.get()) { var records = consumer.poll(Duration.ofSeconds(1)); for (var partition : records.partitions()) { var batch = records.records(partition); for (var record : batch) { var eventId = header(record, "event-id"); if (!inbox.tryInsert(eventId, record.topic(), record.partition(), record.offset())) { continue; } var workflow = client.newWorkflowStub(ModelFlow.class, record.key()); workflow.onEvent(ModelEvent.from(record)); inbox.markApplied(eventId); } var nextOffset = batch.get(batch.size() - 1).offset() + 1; consumer.commitSync(Map.of(partition, new OffsetAndMetadata(nextOffset))); } if (inbox.backlog() > 50_000) consumer.pause(consumer.assignment()); else consumer.resume(consumer.assignment()); } } The durable inbox is the effective-once bridge. If the process crashes after signaling Temporal but before committing offsets, Kafka may redeliver, yet tryInsert suppresses reapplication. If upstream producers use Kafka transactions, the consumer should read with isolation.level=read_committed so aborted records stay invisible; Kafka’s configuration reference notes that read_committed returns only committed transactional messages and withholds records past the last stable offset while open transactions exist. Backpressure also belongs here. Kafka exposes pause and resume without forcing a group rebalance, and monitoring guidance explicitly recommends watching lag, fetch rate, poll timing, and commit latency to ensure consumers are keeping up. Context propagation is easiest when context is split into stable metadata and mutable conversational state. Stable identifiers such as trace ID, tenant, policy version, and conversation key belong in Kafka headers and Temporal headers so they survive hops across services and activities; Kafka’s ProducerRecord supports headers, and Temporal context propagators move custom key-value data across workflow, activity, and child-workflow boundaries. Mutable context, by contrast, should not live in worker memory or ad hoc caches. It belongs in the workflow state, often as a compact summary plus references to offloaded artifacts. Temporal’s documentation explicitly warns that all activity inputs and outputs are persisted, that long AI-style conversations grow history, and that large histories degrade workflow-task latency. For long-running sessions, Continue-As-New provides a checkpoint boundary with a fresh Event History while preserving the workflow identity chain. Reconciliation closes the last reliability gap. Even with careful commits, outages, manual replays, or producer bugs can create suspicion that a workflow missed an event. Temporal queries are read-only and must not mutate state or block, which makes them ideal for asking a workflow for its durable high-water mark and replaying any gap from the event store. Java void reconcile(String workflowId, long durableHighWatermark) { var workflow = client.newWorkflowStub(ModelFlow.class, workflowId); long applied = workflow.lastAppliedSequence(); eventStore.readRange(workflowId, applied + 1, durableHighWatermark) .forEach(workflow::onEvent); } This pattern works because the workflow does not trust delivery history alone; it trusts its own durable state. Observability then becomes the enforcement layer for those guarantees. Kafka should surface lag, request latency, retry rates, poll gaps, and buffer exhaustion, while Temporal should emit metrics through Micrometer, trace activity and workflow paths, and expose searchable workflow metadata through Search Attributes. Temporal also recommends monitoring replay latency because large histories, payload sizes, and cache churn drive recovery cost. Together, these signals reveal the difference between a system that is slow, a system that is duplicating work, and a system that is actually losing context. Conclusion Orchestrating small language models without losing events or context is fundamentally a durability problem, not a prompt-engineering problem. Kafka should be used for ordered transport and scalable ingestion, but semantic completion should be anchored in Temporal’s durable workflow state, where signals, sequence gates, retryable activities, queries, and replay make failures recoverable rather than ambiguous. Exactly-once remains valuable inside Kafka’s transactional envelope, yet end-to-end correctness across model calls and other side effects comes from explicit idempotency, deduplication, reconciliation, and bounded context management with external storage and continue-as-new. In a Java stack, that combination yields an architecture where duplicates become harmless, ordering becomes explicit, back pressure becomes controlled, and context survives crashes because it is recorded in durable state instead of being left in process memory.
VP of Engineering,
Factorial
Founder,
DataView
Vice President - Banking and Finance / Cloud /Bigdata / Analytics / AI & ML,
JPMorgan Chase & Co.
Data Analytics - Assistant Vice President,
U.S. Bank