Integration / ETL / Batch Processing
ETLInternalIntegration Pattern

Integration ETL / Batch Processing (Internal)

David TirabassiUpdated

Problem

An internal backend workload must periodically ingest large, infrequently changing datasets from another internal source within the private subnet. Without an efficient bulk transfer mechanism, ad-hoc or row-by-row movement overloads systems, breaks data consistency, and fails to scale as source volumes grow.

Solution

Implement a data pipeline utilizing a Managed Data Integration Service to automate the extraction, transformation, and loading (ETL) of bulk datasets between internal backend workloads or data stores. The pipeline should be orchestrated to operate on a scheduled basis or triggered by data availability events, ensuring data consistency and integrity across the enterprise data landscape. The service can leverage various deployment models, including serverless functions or containerized workloads, for scalable and cost-effective processing within the Private Subnet (Workloads).

Cloud Paradigm

  • Data Pipelining
  • Batch Processing
  • Serverless Compute (for transformation tasks)
  • Event-Driven Architectures (for triggering)
  • Data Lake / Data Mesh integration readiness
  • Idempotent Data Ingestion
  • Cost-Optimized Data Movement

Solution Flow

Data Ingestion Flow:

  1. Source Backend Workload/Data Store: The source system (e.g., a database, object storage, file share) provides data for extraction. Data can be exposed via secure data interfaces (e.g., database connections, API endpoints, secure file transfer protocols) within a Private Subnet (Workloads).
  2. Managed Data Integration Service: A scheduled trigger or event (e.g., new file arrival, time-based schedule) initiates a job within the Managed Data Integration Service, which operates within a Private Subnet (Workloads).
  3. Data Extraction: The service securely connects to the source, extracts the bulk data, and typically stages it temporarily in an intermediate data store (e.g., object storage or a temporary database).
  4. Data Transformation (Optional/Conditional): The extracted data undergoes transformation processes as required. This might include data cleansing, enrichment, format conversion, aggregation, or joining with other datasets. Serverless functions or containerized compute instances can perform these transformations.
  5. Data Loading: The transformed data is then securely loaded into the target backend workload or data store within a Private Subnet (Workloads). This could involve direct database inserts/updates, bulk file uploads to object storage, or message queue publication.
  6. Target Backend Workload/Data Store: The target system receives and stores the ingested data, making it available for consumption.

When to Use

  • Bulk datasets must move periodically between internal backend systems, and the source changes infrequently enough that continuous streaming would be wasteful.
  • The workload tolerates batch latency (minutes to hours) rather than requiring near-real-time propagation.
  • Data requires meaningful transformation — cleansing, enrichment, format conversion, or joins — before it is usable by the target store.
  • Both source and target reside within Private Subnet (Workloads) and can be reached over secure data interfaces without public exposure.
  • Processing volume varies significantly between runs, making serverless or containerized elastic compute more cost-effective than always-on infrastructure.

When NOT to Use

  • Consumers need sub-second or event-by-event freshness — use a streaming or change-data-capture pipeline instead.
  • The interaction is a request/response exchange between services rather than bulk movement — an API or synchronous integration pattern fits better.
  • Data volumes are small and frequently changing, where per-record eventing is simpler and cheaper than scheduled batches.
  • The transfer crosses a trust boundary to an external partner — an externally-facing, gateway-mediated integration pattern is more appropriate.
  • Source and target share the same database engine and a native replication feature already satisfies the requirement.

Trade-offs

  • Efficient, cost-effective bulk movement vs the inherent batch latency that leaves the target stale between scheduled runs.
  • Elastic serverless/containerized compute vs the operational complexity of orchestration, checkpointing, and dead-letter handling for reliability.
  • Centralized transformation and data quality enforcement vs the added pipeline logic and metadata management that must be maintained as code.
  • Decoupling of source and target schemas vs the risk of silent drift when either side changes without pipeline updates.
  • Strong observability and lineage vs the overhead of instrumenting metrics, logs, and cataloging across every stage.

Real-World Example

Consider a pharmaceutical manufacturer that consolidates batch-release records each night from a manufacturing execution database that barely changes during a production shift. A time-based trigger launches a job in the Managed Data Integration Service within the Private Subnet (Workloads); it connects to the source over TLS, extracts the shift's batch genealogy and quality-test results, and stages them as Parquet files in object storage. A containerized transformation step standardizes lot identifiers, enriches records with reference master data on materials and equipment, and aggregates yield metrics, then bulk-loads the output into the quality-analytics warehouse. Checkpointing and a dead-letter queue capture malformed batch records for reprocessing and audit, while job duration and row counts stream into the central observability platform, giving validation and operations teams the lineage needed for regulatory traceability.

Additional Details

  • Orchestration & Scheduling: Leverage cloud-native schedulers or workflow orchestrators to manage the execution of data pipelines. Jobs can be time-based, event-driven, or manually triggered.
  • Resilience & Error Handling: Design pipelines with built-in fault tolerance, including checkpointing, retry mechanisms, and dead-letter queues for failed records. Implement robust logging and alerting for operational issues.
  • Scalability: Utilize serverless or containerized compute for transformation steps to automatically scale resources based on data volume and processing requirements.
  • Data Governance: Implement data cataloging and metadata management to track data lineage, ensure data quality, and support compliance requirements for the ingested datasets.
  • Data Formats & Protocols: Use efficient data formats suitable for bulk processing (e.g., Parquet, ORC, Avro, CSV, JSON lines). Secure protocols like TLS 1.2+ for all data transfer are mandatory.
  • Observability: Integrate pipeline logs, metrics (e.g., data volume processed, job duration, error counts), and lineage information into a centralized observability platform to monitor performance, health, and data flow.
  • Infrastructure as Code (IaC): Manage the entire data pipeline, including connectivity, compute resources, and orchestration logic, as code for version control, automation, and consistent deployments.

Security Controls

  • Network Segmentation: Deploy data integration components and data stores within isolated Private Subnets (Workloads) with no direct Public Internet access. Restrict network access to the absolute minimum necessary ports and protocols between source, data integration service, and target data stores.
  • Transport Security: Enforce strict Transport Layer Security (TLS 1.2 or higher) for all data in transit between source, processing components, and target.
  • Data Encryption at Rest: Ensure all data stores, including transient staging areas, encrypt data at rest using platform-managed or customer-managed encryption keys.
  • Authentication & Authorization:
    • Utilize managed identity or service accounts with the principle of least privilege for the data integration service to access source and target data stores.
    • Implement role-based access control (RBAC) to manage permissions for pipeline operators and data consumers.
  • Credential Management: Never store credentials in cleartext. Use a Managed Secrets Management Service to store and retrieve database credentials, API keys, or other sensitive access tokens securely.
  • Data Validation & Integrity: Implement robust data validation checks as part of the transformation process to prevent corrupted or malicious data from entering target systems.
  • Audit Logging: Enable comprehensive logging for all data pipeline activities, including data access, transformations, and load events, integrating with a centralized logging and monitoring solution.
  • Data Loss Prevention (DLP): For highly sensitive data, implement DLP scans within the transformation pipeline to identify and prevent unauthorized data egress or inappropriate storage.

Related Patterns