Salesforce bulk data ingestion using Application Integration

By Google Cloud Tech

Share:

Key Concepts

  • Large-scale data ingestion: Efficiently transferring vast amounts of data.
  • Application integration: Connecting different software applications to share data.
  • Salesforce: A cloud-based customer relationship management (CRM) platform, used as the data source.
  • Bulk API v2: Salesforce's optimized API for transferring large datasets, crucial for performance.
  • Asynchronous orchestrator-worker model: An architectural pattern for managing complex, long-running tasks, where an orchestrator coordinates multiple workers.
  • Orchestrator: The component responsible for initiating, monitoring, and managing the overall data ingestion job, including status checks and worker coordination.
  • Worker: The component responsible for the actual data transfer, retrieving data in manageable chunks and uploading them.
  • Google Cloud Storage (GCS): The cloud-based object storage service used as the immediate destination for ingested data.
  • SOQL query (Salesforce Object Query Language): Salesforce's proprietary query language used to define the data and fields to retrieve.
  • Integration Connector Service: A service used to create and manage connections to various data sources, including Salesforce.
  • initiate job API trigger: The specific API call that starts a full data synchronization job.
  • Workload input parameters: Configuration settings for the ingestion job, such as SOQL query, Page size, GCS bucket, Job unique identifier, and Timer.
  • Pagination: The process of dividing a large dataset into smaller, sequential pages for easier processing.
  • Sanity check: A basic verification process to ensure the consistency and correctness of output files.
  • BigQuery: Google Cloud's fully managed, serverless data warehouse, mentioned as a potential final destination for the ingested data.

Overview of Large-Scale Data Ingestion with Application Integration

This demonstration outlines a method for ingesting large datasets at scale using application integration, specifically leveraging Salesforce as a data source and its Bulk API v2 for optimized data transfer. The underlying pattern, however, is adaptable to any connector that supports a bulk or batch processing mechanism.

Asynchronous Orchestrator-Worker Model

The core of this solution is an asynchronous orchestrator-worker model, designed for efficient transfer of large data sets.

  • Initiate Job Task: The process begins with an initiate job task responsible for setting up and starting the ingestion job.
  • Orchestrator's Role: The orchestrator is central to managing the overall progress and job status. Crucially, it triggers the worker processes only after the underlying Salesforce jobs have completed, ensuring proper sequencing and resource management. It also tracks the overall progress.
  • Worker's Role: The worker handles the actual data transfer. It retrieves data from Salesforce in manageable chunks and uploads these chunks to Google Cloud Storage (GCS). This process is recursive and sequential, with the worker calling back the orchestrator to request the next page of data for ingestion, facilitating pagination.

Salesforce Connection and Configuration

To implement this, specific setup steps are required:

  1. Data Preparation: Ensure your Salesforce environment contains the data intended for ingestion. A SOQL query is used to define the specific data and fields to retrieve.
  2. Connection Creation: A new connection must be created within the Integration Connector Service. Select Salesforce and configure the connection settings.
  3. Bulk API v2 Configuration: The Salesforce connection must be specifically configured to work with the Bulk API v2, which is available as part of the standard actions provided by the connector.
  4. Template Utilization: It is highly recommended to start from the provided template titled "large-scale data ingestion Salesforce with bulk API." This template pre-contains the described orchestrator-worker logic and is specifically designed for Salesforce with Bulk API v2.

Initiating the Ingestion Job (Full Sync)

For a full data synchronization, the process is initiated using the initiate job API trigger. This trigger starts the Salesforce Bulk API query job and simultaneously launches the monitoring flow managed by the orchestrator. The job parameters are defined through the workload input, which includes:

  • SOQL query: The Salesforce Object Query Language query specifying the data to retrieve.
  • Page size: Sets the maximum number of records to be processed in a single batch.
  • Next page: A placeholder for pagination, which should be empty at the start of the job.
  • Job status: Starts as "open."
  • GCS bucket and GCS folder: Specify the destination folder within Google Cloud Storage.
  • Job unique identifier: A unique ID required for tracking the job.
  • Timer: Sets the polling interval for checking the job status.
  • Email for notification: An email address for receiving job notifications.

Output Verification and Further Steps

Upon completion of the ingestion job, the results can be verified in the specified GCS folder.

  • File Nomenclature: It is advisable to modify the nomenclature of output files by adding timestamps, a page ID, or even the job ID created for the operation. This practice significantly aids in debugging in case of failures.
  • Sanity Checks: As a crucial next step, users are encouraged to create a sanity check to verify the consistency of the output files, ensuring data integrity.
  • BigQuery Integration: If the final destination for this data is BigQuery, the raw data can be easily transferred into the destination table using BigQuery jobs.

Synthesis and Conclusion

This solution provides a robust and scalable framework for ingesting large datasets from Salesforce using its Bulk API v2, leveraging an asynchronous orchestrator-worker model. By providing a structured approach from connection setup to job initiation and output verification, including a pre-built template, it streamlines the complex process of large-scale data transfer. The emphasis on specific configurations, detailed job parameters, and post-ingestion checks (like file naming for debugging and sanity checks) ensures both efficiency and reliability, making it an actionable blueprint for high-volume data integration.

Chat with this Video

AI-Powered

Load the transcript when you're ready to chat so the initial page stays lighter.

Ready to summarize another video?

Summarize YouTube Video