Wednesday
Hi Everyone, One of our big ingestion tasks is lakeflow connect CDC ingestion of ERP data from a SQL server database. I am deploying the pipelines and jobs for this via DABs. We are ingesting about 80 tables from the database, a handful of which are huge (100s of GBs). These tables take a long time for the initial ingestion, one takes about 36 hours. After that the "steady state" is fine.
I've tried every DAB configuration I can think with min and max workers but I can't get the gateway compute to be elastic, so that it reduces this initial ingestion time by using more workers, then goes down to one or zero workers when that's over. Is this possible? The gateway compute is continuous and it seems to get stuck at one driver and one worker and never expands/contracts.
Wednesday
I ran into a similar issue with SQL Server Lakeflow Connect while ingesting some very large tables. After trying a few different options and discussing it with Databricks, this is what has been explained.
The initial-load parallelism is driven more by the connector internals than by normal Spark autoscaling, for example:
- maximum concurrent connections to SQL Server
- whether the large table can be split into parallel chunks using a suitable split key
- maximum concurrent jobs/chunks
If the split key is not suitable for partitioning, a very large table can effectively be read as one large chunk, so adding workers does not help much. Some of those settings are backend-controlled and may require Databricks Support/account-team involvement.
The public docs also recommend using small gateway workers because worker size itself does not materially improve gateway extraction performance.
One newer option worth evaluating is SQL Server Integrated CDC. It removes the separate continuous gateway entirely and uses one pipeline for extraction + application:
Wednesday
@data_pulse Thanks for your reply. A red flag for me regarding Integrated CDC is the line "Because each extraction stage runs for at least 10 minutes, an interval of 60 minutes or longer is a good starting point".
We have one ERP CDC job that needs to run at 15 minute intervals.
Wednesday
If you need 15 minute CDC interval - Migrate those specific jobs to the integrated CDC pipeline in continuous speed-optimized mode. If the 15 minute jobs involves fewer than 50 tables, it works directly. If not, split them into multiple pipelines. More details here You can let enhanced autoscaling handle the rest.
resources:
pipelines:
continuous_cdc_pipeline:
name: erp_sql_server_pipeline
channel: PREVIEW
continuous: true
catalog: default
schema: sql_server
ingestion_definition:
connection_name: erp_conn
connector_type: CDC
objects:
- table:
source_catalog: ingest
source_schema: ss
source_table: providerFor initial snapshot - investigate whether the source SQL Server can be tuned - try increasing parallelism on the database side, ensuring CDC/change tracking is properly configured and check for lock contention during the snapshot. The gateway's speed is also bounded by how fast SQL Server can serve the data.
Set the mode to ENHANCED explicitly on the gateway auto scale config and check the pipelines again for the 36 hour initial ingestion phase. Split the largest tables into separate gateway pipeline pairs for better ingestion if feasible.
yesterday
Greetings @dbernstein_tp, I did some digging and here is what I found.
Short answer: no, the gateway can't scale out horizontally, and that's by design rather than a DAB problem. The distinction that matters is Spark worker autoscaling versus the connector's own extraction parallelism. data_pulse and balajij8 have already covered good ground, so I'll pull it together.
Why it sits at one driver and one worker. The gateway isn't a normal Spark job that fans the snapshot out across executors. Extraction runs on the driver. The docs say to use the smallest possible worker nodes because they don't affect gateway performance, and the driver needs at least 8 cores for efficient extraction. Their reference compute policy pairs a very large driver (r5n.16xlarge on AWS) with tiny workers (m5n.large). So min_workers and max_workers won't shorten your 36-hour snapshot. balajij8 suggested setting ENHANCED on the gateway autoscale; I'd temper expectations there for the same reason.
The dial you do have is vertical: driver size. The gateway is a pipeline object, so you can pin node types in its clusters block in the bundle, or through a cluster policy and policy_id:
gw_pipeline:
continuous: true
clusters:
- label: default
driver_node_type_id: <8+ core driver, go big for the snapshot>
node_type_id: <smallest worker available in your region>
Two caveats. If the big table has no suitable split key, data_pulse's point stands: the connector may read it as one chunk and a bigger driver only helps so much. Chunking and concurrency are backend-controlled, so that's a Support or account team conversation. And shrinking the driver after the snapshot means a redeploy that restarts the gateway, which the docs warn against because changes can drop if the source truncates logs during the gap. One planned restart inside your CDC retention window is a fair tradeoff against paying for a huge driver forever.
Before you resize anything, check the source side: SQL Server lock contention, the index on the split key, the network path to the gateway, and the pipeline event log. The gateway can't pull faster than SQL Server serves.
The closest thing to horizontal scaling is balajij8's idea of splitting the largest tables into their own gateway and pipeline pairs. Each gateway has its own driver, so three gateways snapshotting at once is three times the extraction capacity. The bill for three continuous classic clusters doesn't go away after the snapshot, though.
Integrated CDC and your 15-minute requirement. The 10-minute minimum and 60-minute starting interval you flagged apply to triggered mode. Continuous mode (Beta, enabled from the Previews page) is the fit for a hard 15-minute target. Speed-optimized mode (pipelines.managedIngestion.continuous.runMode: SPEED) is documented as applying changes within a few minutes with aggressive autoscaling, capped at 50 tables per pipeline, so your 80 tables need at least two pipelines, as balajij8 noted. The default scale-optimized mode handles up to 500 tables but rotates across them without aggressive autoscaling, and the docs don't put a latency number on it. Both continuous modes put higher load on the source than triggered mode, so confirm that's acceptable.
Before you jump: the integrated connector needs a workspace feature flag from your account team, integrated pipelines get vertical autoscaling by default (an OOM on one update means a bigger driver next time), and a full refresh in continuous mode may need multiple restarts because the snapshot is staged asynchronously. It fixes the "paying for an idle gateway" problem. It doesn't necessarily fix the "36 hours" problem.
Practical next step: put your biggest table or two in a separate test pipeline, record source read throughput, driver memory, and end-to-end lag, and compare the gateway design against integrated CDC before committing. Also confirm the source is a primary instance (replicas aren't supported) and that your log retention window is longer than your worst-case catch-up time.
Takeaway: make the driver big for the snapshot, split the monster tables across gateways if the budget allows, get Support involved if the split key is the real bottleneck, and benchmark continuous speed-optimized integrated CDC for the 15-minute workload.
References:
clusters block: https://docs.databricks.com/aws/en/ingestion/lakeflow-connect/sql-server-pipelineRegards, Louis.
yesterday
Great discussion. I think there is another architectural consideration worth exploring before changing the gateway configuration.
Does the 15-minute freshness requirement apply to all 80 tables, or only to a subset of business-critical tables?
If only a subset requires low latency, I would consider separating the ingestion strategy into different workload groups:
Business-critical tables with strict freshness requirements, potentially using continuous speed-optimized Integrated CDC.
Less time-sensitive tables using scheduled ingestion or a more cost-efficient mode.
Very large tables evaluated separately during their initial snapshot to avoid unnecessary contention with ongoing production workloads.
I would also validate whether Change Tracking is sufficient for tables with primary keys and where only the latest state is required, since it can reduce the impact on SQL Server. Tables requiring complete change history would need a different approach.
Finally, I would measure actual end-to-end data freshness rather than relying exclusively on pipeline execution intervals. A pipeline running every 15 minutes does not necessarily guarantee that the data is available within 15 minutes, especially during backfills or recovery scenarios.
Is the 15-minute requirement shared by all tables, or could the workload be separated according to business SLAs?
yesterday
About 21 of the tables are in the 15 minute refresh job. These are cleaned and then aggregated into six gold zone tables, two of which are used in dashboard reporting. The other 60 tables and their aggregations are in a job that is on a one hour schedule.
I played with continuous ingestion earlier this year and we had to turn it off because it was too expensive. I have not experimented yet with the integrated CDC feature.
yesterday
Thanks for clarifying! It sounds like you've already made a good architectural decision by separating workloads based on their refresh requirements.
Since continuous ingestion was too expensive, I wouldn't recommend switching approaches without first identifying the actual bottleneck.
I'd focus on measuring three stages separately: source extraction, Bronze ingestion, and downstream Silver/Gold transformations.
Since only two Gold tables feed the dashboards, it might also be worth identifying which of the 21 source tables are actually on their critical path.
Integrated CDC could be interesting to evaluate with a small subset of tables, comparing end-to-end latency, source impact, and cost before considering a broader migration.
One question: is your biggest challenge currently the 36-hour initial historical load, or are you also struggling to consistently meet the 15-minute refresh requirement during incremental ingestion?