Data Ingestion and Source System Integration Questions
Getting data out of heterogeneous source systems and landing it reliably: APIs, operational databases, file drops, webhooks, message queues and third-party SaaS. Covers connector selection and design (managed platforms versus Debezium, DMS or Kafka Connect versus building your own), pull versus push and polling versus webhook patterns, incremental extraction and high-watermark strategy including what to do when a source offers no native change capture, authentication and credential rotation against third-party APIs, source-side rate limits and quotas, schema drift and contract breakage at the source boundary, backfill and replay of history, ingestion-time data-quality gates, reconciliation after a source outage, and negotiating with source-system owners. The scope stops at the boundary: once data has landed, transforming it, the architecture of the pipeline that carries it, stream-processing mechanics, and pipeline monitoring are all covered separately.
Describe how you would secure ingestion pipelines that span multiple cloud services and on-prem sources. Cover authentication and authorization for connectors (mTLS, IAM roles, service accounts), encryption at rest and in transit, secret management, auditing, and how you enforce least privilege for producers and consumers on both sides of a connector.
Sample Answer
Direct answer
Securing an ingestion pipeline that spans multiple clouds and on-prem sources means treating every connector as a distinct trust boundary: authenticate it with the strongest mechanism the source and target both support, encrypt data both in transit and at rest, keep every secret in a managed secrets store with rotation, log enough to reconstruct who touched what and when, and grant each connector only the specific permissions it needs, never a broad standing credential.
Structured elaboration
Authentication and authorization
- Prefer mutual TLS (mTLS) for connector-to-connector traffic where both ends support it, since it authenticates both parties, not just the client to the server.
- Use cloud-native IAM (Identity and Access Management) roles or service accounts for connectors talking to a cloud provider's own services, rather than long-lived static keys; a role can be scoped narrowly and its usage is natively auditable.
- For on-prem sources with neither mTLS nor cloud IAM available, a short-lived, narrowly-scoped service account credential is the fallback, still avoiding a single shared "integration user" account used by every connector.
Encryption
- In transit: TLS for every hop, including internal ones between a connector and an internal message bus; do not assume "internal network" means encryption is unnecessary.
- At rest: envelope encryption via a key management service (KMS), so data is encrypted with a data key that is itself encrypted by a master key the KMS controls, giving you centralized key rotation and revocation without re-encrypting all your data on every rotation.
Secret management
- Every credential, API key, database password, and certificate lives in a secrets manager, never in connector configuration files, environment variables checked into source control, or container images.
- Automate rotation where the source supports it; where it does not, at minimum track credential age and alert on staleness so rotation happens on a schedule rather than never.
Auditing
- Log every authentication event, every credential access, and every connector's data-access pattern (which source, how much data, when) in a way that is queryable after the fact, not just written to an ephemeral log that rotates away in days.
- Treat audit logs as themselves sensitive: they can reveal a lot about your data flows and access patterns, so restrict who can read them.
Least privilege for producers and consumers
- A connector reading from a source should have read-only access to exactly the tables, topics, or endpoints it needs, never blanket access to the whole source system.
- A connector writing to a target should be scoped to write only to its designated landing location, not broad write access across the destination platform.
- Review these scopes periodically; connector permissions tend to accumulate over time as requirements evolve, and nobody circles back to remove access that is no longer needed.
Worked example
A connector reads from an on-prem PostgreSQL database (via a VPN tunnel with TLS) and writes to a cloud data warehouse. The read side authenticates with a database service account scoped to SELECT on exactly the three tables the connector needs, credentials stored in the cloud secrets manager and fetched at connector startup rather than baked into its container image. The write side authenticates via a cloud IAM role scoped to write only into the specific landing schema this connector owns, with no access to any other team's tables. Both the credential fetch and every batch write are logged with a timestamp, source, and row count, so a security review six months later can answer "what did this connector actually access and when" without needing to reconstruct it from memory.
Trade-offs & pitfalls
- The most common real-world shortcut is a single shared "integration" database user with broad access used by every connector "to keep things simple"; this collapses your entire audit trail (you cannot tell which connector did what) and turns one leaked credential into access to everything.
- mTLS is not always available on legacy on-prem systems; do not let its absence become an excuse to skip TLS entirely, a one-way TLS connection is still meaningfully better than plaintext.
- Rotating secrets automatically is only safe if the connector can pick up a rotated credential without a manual restart; verify this before relying on automated rotation, or a rotation event becomes an unplanned outage.
- Least-privilege scopes decay over time as requirements shift; without a periodic review, "least privilege at launch" quietly becomes "broad privilege nobody remembers granting" within a year or two.
Compare Kafka Connect, AWS DMS, and Debezium as connector technologies for moving data out of a source system. For each, discuss the sources and targets it supports, its operational model, latency characteristics, and how it handles schema changes, and name a scenario where you would prefer each one.
Sample Answer
Direct answer
Kafka Connect is a connector framework, not a connector itself: it gives you a plugin runtime, offset management, and a scaling model, but you still choose a source connector (JDBC, Debezium, or a vendor plugin) to run inside it. Debezium is a log-based change-data-capture connector, most often deployed as a Kafka Connect source plugin, that reads a database's transaction log directly. AWS DMS (Database Migration Service) is a managed, standalone replication service that does its own log reading and its own delivery, with no dependency on Kafka at all. The practical choice usually comes down to whether you are already committed to Kafka as your transport, how much operational ownership you want, and whether the source and target are both natively supported by a managed service.
Structured elaboration
Kafka Connect (the framework)
- Sources and targets: whatever plugin you install; huge open-source and commercial ecosystem (JDBC, Debezium, S3, Elasticsearch, Snowflake sinks, and more).
- Operational model: you run and scale Connect workers yourself (or use a managed flavor like Confluent Cloud or MSK Connect); connectors are configured via a REST API and distribute their work as tasks across workers.
- Latency: near-real-time when paired with a log-based source connector like Debezium; can be batch-like (minutes) with a JDBC polling connector.
- Schema handling: integrates with a schema registry; Single Message Transforms (SMTs) let you reshape, mask, or route records in-flight without writing custom code, for example dropping a field or renaming a topic based on the source table.
Debezium (a log-based CDC connector, usually run inside Kafka Connect)
- Sources: MySQL, PostgreSQL, MongoDB, SQL Server, Oracle, and others, reading each database's native change log (binlog, write-ahead log or WAL, oplog).
- Targets: Kafka topics (one per captured table by default); further delivery from there is a separate sink connector's job.
- Operational model: you run it, typically as a Kafka Connect source connector; it needs log-retention and permissions on the source database, and captures deletes and DDL, not just row values.
- Latency: sub-second to low-second, because it tails the log rather than polling.
- Schema handling: emits schema-aware records (with before/after state) and can be paired with a schema registry; DDL changes on the source generate corresponding schema evolution events.
AWS DMS (a managed, standalone replication service)
- Sources and targets: a fixed, broad matrix of relational and NoSQL engines on both ends (Oracle, SQL Server, PostgreSQL, MySQL, DynamoDB, S3, Redshift, and more), configured declaratively rather than through a plugin ecosystem.
- Operational model: fully managed by AWS; you configure a replication instance and endpoints and it handles the full-load-plus-CDC lifecycle, including the initial snapshot, without you standing up any infrastructure.
- Latency: log-based CDC once the initial load completes, comparable to Debezium; the initial full load is a separate, often slower phase.
- Schema handling: less flexible in-flight transformation than Kafka Connect SMTs; schema changes on the source can require re-running or reconfiguring a task rather than being absorbed automatically.
When you would prefer each
- Prefer Kafka Connect + Debezium when Kafka is already your transport of record and you want every downstream consumer, not just one target, to see the change stream, or when you need in-flight transformation via SMTs.
- Prefer AWS DMS when you are moving data between two AWS-native or classic relational engines and do not want to operate any connector infrastructure yourself, especially for a one-time or few-target migration rather than a fan-out to many consumers.
- Prefer Kafka Connect with a non-CDC plugin (JDBC source, a vendor-published connector, a sink connector) when the source has no usable log access at all and you are accepting timestamp-based polling instead.
Worked example
A team is migrating change events from an on-prem PostgreSQL order-processing database so three different services (fraud detection, analytics, and a search index) can each consume the same stream independently. Because there are three independent downstream consumers, not one target, Kafka Connect with Debezium is the right shape: Debezium publishes one topic per table, and each service runs its own consumer group against those topics with no coupling to the others. If instead the requirement had been "replicate this same database into an RDS PostgreSQL replica for read scaling, one source, one target," AWS DMS would have been the simpler, lower-operational-cost choice, since there is nothing to fan out and no need to run Connect infrastructure at all.
Trade-offs & pitfalls
- Debezium needs enough log retention on the source (binlog expiry, WAL retention) to survive a connector outage without falling behind and being forced into a full re-snapshot; this is an easy production surprise if nobody sizes it.
- Kafka Connect's flexibility is also its operational cost: you own worker sizing, connector upgrades, and dead-letter handling for bad records, none of which AWS DMS asks of you.
- AWS DMS's schema-change story is the sharpest limitation: adding or renaming a column mid-flight is not something it absorbs gracefully the way a schema-registry-backed Debezium pipeline can.
- Do not conflate "Kafka Connect" and "Debezium": Connect is the runtime, Debezium is one of many possible connectors you can run inside it, and mixing them up in an interview answer signals you have not actually operated either.
Design the integration layer for a company with 25 internal applications and 8 external vendors that all share customer, product, and order data. Requirements: one place to manage ownership, support for both near-real-time and batch sync, lineage you can use for audits, and minimizing point-to-point connections between systems. How would you structure the data flow and the integration boundaries?
Sample Answer
Direct answer
The way to avoid 25 x 8 point-to-point connections is a hub-shaped integration layer: every internal application and external vendor connects to ONE shared layer, not to each other directly, and that shared layer owns the canonical model, the sync mechanics, and the lineage, while individual systems only need to know how to talk to the hub.
Structured elaboration
Structuring the data flow
- A central integration layer (built on an event bus plus a set of connectors, or an iPaaS-style platform (integration platform as a service: a vendor-hosted product that ships pre-built connectors and a workflow engine so you configure integrations rather than code them from scratch), depending on scale and budget) sits between every system; a system publishes changes to the hub and consumes what it needs from the hub, never directly from another system's database or API.
- Define a canonical data model for the shared entities (customer, product, order) that every system's data is mapped into on the way in and mapped out of on the way out, so the hub, not each pairwise connection, absorbs the complexity of reconciling different systems' representations of the same entity.
Ownership
- Each canonical entity has one clearly designated owning system (the source of truth for that entity), with every other system treated as a consumer or, where genuinely necessary, a secondary contributor with an explicit conflict-resolution rule.
- The integration layer itself owns the mapping rules and the sync mechanics, but not the business data; ownership of "what is correct" stays with the source systems, ownership of "how it moves" belongs to the integration layer.
Supporting both near-real-time and batch sync
- Systems that can emit events (most modern internal apps, some vendors) integrate via the event bus for near-real-time propagation.
- Systems that cannot (older internal apps, vendors offering only file exports or low-frequency APIs) integrate via scheduled batch connectors into the same hub, landing into the same canonical model on their own cadence.
- The consuming side does not need to know or care which mechanism produced an update; it reads from the hub's canonical store either way.
Lineage for audits
- Every record moving through the hub carries provenance: which source system it came from, when it was captured, and which mapping rule transformed it into the canonical shape.
- Centralizing this in the hub, rather than expecting each of the 33 systems to independently implement its own audit logging, is what makes a company-wide audit actually answerable in one place instead of 33 separate investigations.
Minimizing point-to-point connections
- Every pair of systems that needs to exchange data would need its own direct connection without a hub; the number of possible pairs among N systems is written C(N,2) (read "N choose 2") and works out to N x (N-1) / 2, since each of the N systems can pair with N-1 others and dividing by 2 removes the double-count of listing each pair twice (once from each system's own perspective). Without a hub, the theoretical maximum of pairwise integrations among all 33 systems, since the question states all three data types are shared broadly rather than only flowing from vendor to internal app, is C(33,2) = 33 x 32 / 2 = 528, not just the 25 x 8 = 200 cross-group (internal-to-vendor) pairs alone: internal-to-internal sharing among the 25 apps adds up to C(25,2) = 300 more possible pairs, and vendor-to-vendor sharing among the 8 vendors adds up to C(8,2) = 28 more. The hub collapses all of that, cross-group and within-group alike, down to 33 connections total (one per system, in and out of the hub), which is the entire point of centralizing rather than letting teams integrate pairwise as needs arise.
- New systems joining later (a 26th internal app, a 9th vendor) need only ONE new connection to the hub, not a new connection to every existing system they might need data from.
Worked example
A new CRM vendor is being onboarded as external vendor #9. Without the hub, onboarding it would mean building point-to-point integrations to however many of the 33 other systems need customer data from it, potentially dozens of new connections. With the hub, the CRM vendor gets exactly one connector mapping its customer records into the canonical customer model, publishing changes to the hub; every one of the other 32 systems that already consumes canonical customer data from the hub automatically sees the new vendor's data with zero changes on their end, since they were never coupled to any specific source system in the first place, only to the hub's canonical shape.
Trade-offs & pitfalls
- A canonical data model is genuine, ongoing modeling work, not a one-time exercise; as new systems join with data that does not map cleanly onto the existing canonical shape, the model itself has to evolve, and that evolution needs the same schema-change discipline as any other shared contract.
- Centralizing everything into one hub creates a single dependency every system now relies on; the hub's own availability and scaling become critical-path concerns in a way that no single point-to-point connection ever was on its own.
- Not every legacy vendor will cleanly fit the hub pattern; some very old or inflexible systems may need a bespoke adapter in front of them just to get their data into a shape the hub can accept at all, and pretending otherwise just pushes the complexity into a fragile, undocumented corner.
- "One connection per system" is the right target, but it is easy to accidentally recreate point-to-point sprawl inside the hub itself if teams start building direct, hub-bypassing shortcuts "just this once" for a deadline; governance over the hub's own boundary matters as much as the initial architecture.
You need to backfill two years of historical, paginated data for many accounts from a third-party API that is capped at 10 requests per second, while a live incremental sync keeps running against the same API. Describe your parallelization strategy, how you checkpoint so the backfill can resume, how you coordinate the rate limit across workers, and how you guarantee the result is eventually consistent without duplicating records.
Sample Answer
Direct answer
The core tension is that the backfill and the live sync are both consuming the same 10 requests/second budget from the same API, so the design has to explicitly partition that budget rather than let the two compete unpredictably, while making the backfill itself resumable and idempotent so a multi-day job surviving several restarts still converges on a correct, non-duplicated result.
Structured elaboration
Partitioning the rate-limit budget
- Reserve a fixed share of the 10 req/s for live sync (enough to keep it comfortably within its freshness service-level agreement (SLA)) and let the backfill use the remainder, rather than a naive "whoever gets there first" free-for-all between the two workloads.
- A shared, centralized rate limiter (not one limiter per worker) is required here: independent per-worker limiters cannot see each other's usage and will collectively exceed the source's real cap.
Parallelization strategy
- Parallelize across accounts, not within a single account's page sequence, since pages within one account's history typically must be walked in order for correct checkpointing, while different accounts are fully independent of each other.
- Size the worker pool to the rate-limit budget available to the backfill, not to raw compute capacity; more workers than the rate limit can support just means more of them sitting idle waiting for a token.
Checkpointing for resumability
- Checkpoint per account, independently: persist each account's furthest-completed page or cursor so a restart resumes only the accounts still in progress, not the ones already fully backfilled.
- Persist checkpoints frequently enough (after every page, not just at the end of an account) that a crash loses at most one page of already-fetched-but-uncommitted work per in-progress account.
Guaranteeing eventual consistency without duplicates
- Every record written by either the backfill or the live sync goes through the same idempotent upsert path, keyed on the record's own stable ID, so it does not matter which of the two processes writes a given record first or whether both happen to write it.
- Define a clear ordering rule for the rare case where both processes touch the same record concurrently (for example, "whichever write carries the later
updated_atwins"), so the result is deterministic rather than a race.
Worked example
2,000 accounts need two years of history backfilled, at a shared 10 req/s budget, with 3 req/s reserved for live sync, leaving 7 req/s for the backfill. A pool of 7 backfill workers, each holding one token from a centralized rate limiter, is assigned accounts from a queue; each worker walks one account's full page history, checkpointing its cursor after every page, before picking up the next account from the queue. If the process crashes after 1,200 of the 2,000 accounts are fully done and a 1,201st is half-complete, a restart re-reads the checkpoint table, skips the 1,200 complete accounts entirely, resumes the 1,201st from its last saved cursor rather than page 1, and continues the remaining 799 from scratch. Throughout, both the backfill and the live sync write through the same upsert-by-record-ID path, so if live sync happens to ingest a very recent change to an account the backfill has not yet reached, the backfill's later arrival at that same record, carrying older data, does not overwrite the newer one, because the upsert conflict rule explicitly prefers the later updated_at.
Trade-offs & pitfalls
- Splitting the rate-limit budget statically (a fixed 3/7 split) is simple but can starve one workload if the actual demand shifts, for example if live sync traffic spikes; a smarter allocator that lets the backfill temporarily borrow idle live-sync capacity is more efficient but meaningfully more complex to build and reason about correctness for.
- Checkpointing at the account level but not the page level within an account means a crash can lose up to one full account's progress, not just one page; the finer-grained the checkpoint, the less repeated work a crash costs, at the price of more frequent writes to the checkpoint store.
- The "later
updated_atwins" conflict rule assumes the source's own timestamps are trustworthy and consistent across accounts and over the two-year backfill window; if that assumption does not hold for some historical data, you need a different, explicit tie-breaking rule rather than silently trusting a timestamp that might be wrong. - A single centralized rate limiter is also a single point of contention; at very high worker counts it can itself become a bottleneck, though at 7 concurrent workers against a 7 req/s budget this is not yet a real concern.
You are integrating three SaaS systems into your warehouse. One emits events, one only supports paginated reads, and one exports a file every night. The business wants a daily dashboard now and near-real-time alerts from one of the sources later. How would you choose the integration pattern for each source, and keep the overall design maintainable as the requirements evolve?
Sample Answer
Direct answer
Match the integration pattern to what each source can actually do, not to a single company-wide standard: the event-emitting source becomes a subscriber, the paginated-read-only source becomes a scheduled poll, and the nightly-file-export source becomes a scheduled file pickup. The design stays maintainable as requirements evolve by landing all three into the same normalized raw layer with a consistent schema and metadata, so a later change (adding near-real-time alerts from one source) only touches that one source's ingestion path, not the shared downstream model.
Structured elaboration
Matching pattern to source capability
- The event-emitting source: subscribe to its events (a webhook or a stream), land each event as it arrives; this is the only one of the three sources that can support the future near-real-time alert requirement without a redesign.
- The paginated-read-only source: a scheduled poll walking every page since the last checkpoint, on whatever interval balances freshness against the source's rate limit.
- The nightly-file-export source: a scheduled pickup job (SFTP, meaning SSH File Transfer Protocol, cloud storage, or however the export is delivered) that runs shortly after the export is expected, with a check that the file has actually landed before processing it.
Keeping the design maintainable
- Land all three into a common raw/landing layer with a consistent envelope (source name, ingestion timestamp, and the original payload), so downstream consumers work against one shape regardless of which pattern produced the data.
- Keep each source's connector isolated: a change to how the file-export source is polled should never require touching the event-subscriber's code.
- Make the freshness characteristics of each source visible downstream (a "last updated" or "as of" indicator per source) rather than presenting all three as if they were equally fresh, since silently blending a nightly-batch source with a real-time one can mislead a dashboard's audience about how current the numbers actually are.
Building in room for the requirement to evolve
- Today's ask is a daily dashboard, which the paginated and file-export sources already satisfy on their existing schedules; the event source's daily aggregation is just a rollup of its already-landed real-time stream, so no rework is needed there either.
- The stated future need (near-real-time alerts from one source) should specifically inform WHICH source you architect for streaming from day one, even if you are not yet building the alerting logic; retrofitting an already-batch-oriented connector into a streaming one later is a much larger project than building it as a subscriber from the start when you already know that is where it is headed.
Worked example
The event-driven source publishes order-status-change events; land each event onto a raw topic or table immediately as it arrives, since this is the source flagged for a future near-real-time alert. The paginated-read-only source (a support-ticketing API) is polled every 30 minutes using cursor-based pagination and a persisted checkpoint, comfortably supporting the daily dashboard's freshness need without over-polling a source that has no true real-time signal to offer anyway. The nightly-file-export source (a partner's inventory feed) is picked up by a job scheduled 30 minutes after the partner's documented export window, with a presence check before processing so a late or missing file triggers an alert rather than the job silently processing yesterday's stale file again. All three land into the same raw schema (source, ingested_at, payload), and the daily dashboard is built from a view over all three, each contributing at its own natural freshness. When the near-real-time alert requirement actually arrives, only the event source's already-real-time data needs a new consumer built against it: neither the poll-based nor the file-based ingestion paths need to change at all.
Trade-offs & pitfalls
- The most common mistake is trying to force all three sources into ONE integration pattern for "consistency," typically by wrapping the event source in a batch poll of its own API instead of subscribing directly, which throws away the one capability (real-time delivery) that source actually offers and that the stated future requirement will need.
- A file-pickup job with no presence check is a classic silent-staleness bug: if the partner's export is late or fails, a naive job that just re-reads "the file at this path" will happily reprocess yesterday's data with no error at all.
- Landing all three sources into a shared raw schema is good for maintainability, but resist the temptation to also force them into the SAME refresh cadence in that shared layer; an artificial "we refresh everything hourly" schedule either wastes effort re-polling a nightly file source or under-serves the event source's real freshness.
- Do not let "keep it maintainable" become an excuse to over-engineer a generic pluggable-connector framework for just three sources; isolating each connector's code is enough, a fully abstracted framework is a cost worth paying only once you have many more sources than three.
Unlock Full Question Bank
Get access to all 12 Data Ingestion and Source System Integration interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.