如何通过Azure Data Factory v2实现本地数据库数据转JSON并流转至Event Hub?
Hey there! Since you're new to Azure Data Factory, let's walk through a clear, step-by-step solution to meet your requirement: moving data from on-premise Oracle and SQL Server to Blob Storage (with one JSON file per row), then routing those files to Event Hub. I'll break this down into manageable parts and share key tips along the way.
Before diving into pipelines, you need to lay the groundwork to connect your local databases to ADF:
- Deploy a Self-Hosted Integration Runtime (SHIR): Since your Oracle and SQL Server are on-premise, ADF can't reach them directly. Install SHIR on a local machine (or a VM with access to your databases), then register it to your ADF instance. This acts as the bridge between ADF and your local data sources.
- Grant Permissions: Ensure the SHIR service account has read access to your Oracle/SQL Server databases. Also, make sure your ADF managed identity (or service principal) has Storage Blob Data Contributor access to your target Blob Storage, and Azure Event Hubs Data Sender access to your Event Hub.
Next, define connections to all your data sources and sinks:
Linked Services
- Oracle/SQL Server Linked Service: Create one for each database, select the SHIR you deployed as the integration runtime, and enter your database credentials.
- Blob Storage Linked Service: Connect to your Azure Storage account (use managed identity for authentication if possible).
- Event Hub Linked Service: Connect to your Event Hub namespace and target Event Hub, again using managed identity for secure access.
Datasets
- Source Datasets: Create datasets for your Oracle and SQL Server tables (or use query-based sources if you only need a subset of data).
- Blob Sink Dataset: Create a JSON-formatted dataset pointing to your target Blob container/folder. Keep the file path dynamic for later use (we'll generate unique filenames per row).
- Event Hub Sink Dataset: Create a dataset linked to your Event Hub, configured to send event payloads.
You have two main options here, depending on your data volume:
Option A: For Small-to-Medium Data Volumes (Lookup + ForEach + Copy)
This is straightforward for smaller datasets:
- Lookup Activity: Add a Lookup activity to fetch all rows from your source table (use a query like
SELECT * FROM your_target_table). Note: By default, Lookup returns 5000 rows—increase this in the activity settings if you need more, but avoid this for large datasets (100k+ rows). - ForEach Activity: Wire the Lookup output to a ForEach activity, setting
Itemsto@activity('YourLookupActivityName').output.value. You can run this in parallel (adjust concurrency in settings) or sequentially. - Copy Activity Inside ForEach:
- Source: Use an Inline JSON dataset as the source. Set the content to
@string(item())—this converts the current row from the Lookup into a JSON string. - Sink: Use your Blob dataset. Set the file path to something like
your-container/your-folder/@{guid()}.json—theguid()function ensures each file has a unique name, preventing overwrites. In the sink format settings, select Single document to ensure each file contains one JSON object (your row data).
- Source: Use an Inline JSON dataset as the source. Set the content to
Option B: For Large Data Volumes (Data Flow)
Data Flow is optimized for bulk processing and avoids the row limits of Lookup:
- Create a Data Flow: Add a new data flow resource in ADF.
- Source Transformation: Connect to your Oracle/SQL Server dataset, select your table or enter a query to filter data.
- Surrogate Key Transformation: Add this to generate a unique, auto-incrementing ID column (e.g.,
row_idstarting at 1). This lets us partition data into individual files per row. - Sink Transformation: Connect to your Blob dataset, set format to JSON. Key settings:
- Under Settings, set Partition option to Column and select the
row_idcolumn. This creates one file per uniquerow_id(i.e., one file per row). - Under Format, select Single document to ensure each file has a single JSON object.
- Under Settings, set Partition option to Column and select the
Once your JSON files are in Blob Storage, you can automate sending them to Event Hub with one of these methods:
Method 1: ADF Pipeline (Batch Processing)
Use this if you want to run scheduled batches:
- Get Metadata Activity: Point to your Blob folder, select
childItemsin the Field list to fetch all JSON files. - ForEach Activity: Iterate over the file list, using
@activity('YourGetMetadataActivity').output.childItemsas theItems. - Copy Activity Inside ForEach:
- Source: Use your Blob dataset, set the file path to
@item().path. - Sink: Use your Event Hub dataset. The Copy activity will automatically send the entire JSON file content as an event payload.
- Source: Use your Blob dataset, set the file path to
Method 2: Trigger-Based Automation (Real-Time/Quasi-Real-Time)
For near-instant processing when new files are added to Blob Storage:
- Azure Logic Apps: Create a Logic App with a When a blob is added or modified trigger. Add a Get blob content action to read the file, then a Send event action to push the content to Event Hub.
- Azure Functions: Use a Blob Trigger function (in Python, C#, etc.) that triggers when a new JSON file is uploaded. The function reads the file content and sends it to Event Hub via the Event Hubs SDK.
- SHIR Performance: If you're moving large datasets, set up a SHIR cluster (multiple nodes) to handle higher throughput. Monitor the SHIR's CPU/memory usage to avoid bottlenecks.
- Data Validation: Add a Derived Column transformation in Data Flow (or a Data Flow Debug session) to check for invalid data (e.g., special characters in strings, date format issues) before writing to Blob.
- Incremental Loads: If you need to sync data regularly, add a timestamp or incremental ID column to your source tables. Modify your Lookup/Data Flow source query to only fetch rows added/updated since the last pipeline run.
- Monitoring: Enable ADF pipeline monitoring to track run statuses and errors. Use Azure Monitor to set up alerts for failed runs or performance issues.
- Cost Optimization: For large data flows, use Azure Integration Runtime with auto-scaling to reduce costs when not processing data. Avoid unnecessary parallelism in ForEach activities to prevent hitting Blob/Event Hub throttling limits.
内容的提问来源于stack exchange,提问作者Hillol Saha

