Update أو merge records in قاعدة بيانات Azure SQL باستخدام دالات Azure

حاليا، يدعم Azure Stream Analytics (ASA) فقط إدراج (إضافت) الصفوف إلى مخرجات SQL (Azure SQL Databases و Azure Synapse Analytics). تناقش هذه المقالة الحلول البديلة لتمكين UPDATE، UPSERT، أو MERGE على قواعد بيانات SQL باستخدام دالات Azure كطبقة وسيطة.

يتم عرض الخيارات البديلة ل دالات Azure في النهاية.

المتطلبات

يمكنك كتابة البيانات إلى الجدول باستخدام أحد الأوضاع التالية:

وضع عبارة T-SQL المكافئة المتطلبات
إلحاق إدراج بلا
الاستبدال دمج (UPSERT) مفتاح فريد
جمع استخدام MERGE (UPSERT) مع عامل التشغيل operator (+=, -=...) مفتاح فريد ومجمع

لتوضيح الفروقات، فكر فيما يحدث عند تناول السجلين التاليين:

Arrival_Time Device_Id Measure_Value
10:00 و 1
10:05 و 20

في وضع الإضافة ، تدرج سجلين. عبارة T-SQL المكافئة هي:

INSERT INTO [target] VALUES (...);

مما ينتج عنه:

Modified_Time Device_Id Measure_Value
10:00 و 1
10:05 و 20

في وضع الاستبدال ، تحصل فقط على القيمة الأخيرة حسب المفتاح. هنا تستخدم Device_Id كمفتاح. العبارة المكافئة ل T-SQL هي:

MERGE INTO [target] t
USING (VALUES ...) AS v (Modified_Time,Device_Id,Measure_Value)
ON t.Device_Key = v.Device_Id
-- Replace when the key exists
WHEN MATCHED THEN
    UPDATE SET
        t.Modified_Time = v.Modified_Time,
        t.Measure_Value = v.Measure_Value
-- Insert new keys
WHEN NOT MATCHED BY t THEN
    INSERT (Modified_Time,Device_Key,Measure_Value)
    VALUES (v.Modified_Time,v.Device_Id,v.Measure_Value)

مما ينتج عنه:

Modified_Time Device_Key Measure_Value
10:05 و 20

وأخيرا، في وضع التراكم ، تجمع Value بعامل تعيين مركب (+=). هنا أيضا تستخدم Device_Id كمفتاح رئيسي:

MERGE INTO [target] t
USING (VALUES ...) AS v (Modified_Time,Device_Id,Measure_Value)
ON t.Device_Key = v.Device_Id
-- Replace and/or accumulate when the key exists
WHEN MATCHED THEN
    UPDATE SET
        t.Modified_Time = v.Modified_Time,
        t.Measure_Value += v.Measure_Value
-- Insert new keys
WHEN NOT MATCHED BY t THEN
    INSERT (Modified_Time,Device_Key,Measure_Value)
    VALUES (v.Modified_Time,v.Device_Id,v.Measure_Value)

مما ينتج عنه:

Modified_Time Device_Key Measure_Value
10:05 و 21

لاعتبارات الأداء، تدعم محولات إخراج قاعدة بيانات ASA SQL حالياً وضع الإلحاق فقط. تستخدم هذه المحولات إدراجاً ضخماً لزيادة المعدل نقل والحد من الضغط الخلفي.

توضح هذه المقالة كيفية استخدام دالات Azure لتنفيذ وضعي الاستبدال والتراكم لـ ASA. عندما تستخدم دالة كطبقة وسيطة، فإن أداء الكتابة المحتمل لا يؤثر على وظيفة البث. في هذا الصدد، يعمل استخدام دالات Azure بشكل أفضل مع Azure SQL. باستخدام Synapse SQL، قد يؤدي التبديل من عبارات مجمعة إلى عبارات صف تلو صف إلى حدوث مشكلات أداء أكبر.

دالات Azure output

في هذه المهمة، تستبدل مخرج ASA SQL بمخرج ASA دالات Azure. تقوم هذه الدالة بتنفيذ قدرات التحديث أو UPSERT أو الدمج.

حاليا، يمكنك الوصول إلى قاعدة بيانات SQL في دالة باستخدام خيارين. الخيار الأول هو ربط مخرجات Azure SQL. حاليا محدود ب C#، ويقدم فقط وضع الاستبدال. الخيار الثاني هو كتابة استعلام SQL ليتم تقديمه عبر برنامج تشغيل SQL المناسب (Microsoft. Data.SqlClient ل .NET).

تفترض العينتان التاليتان مخطط الجدول التالي. يتطلب خيار الربط مفتاحاً أساسياً ليتم تعيينه في الجدول الهدف. إنه ليس ضرورياً، ولكن يوصى به، عند استخدام برنامج تشغيل SQL.

CREATE TABLE [dbo].[device_updated](
	[DeviceId] [bigint] NOT NULL, -- bigint in ASA
	[Value] [decimal](18, 10) NULL, -- float in ASA
	[Timestamp] [datetime2](7) NULL, -- datetime in ASA
CONSTRAINT [PK_device_updated] PRIMARY KEY CLUSTERED
(
	[DeviceId] ASC
)
);

لاستخدام دالة كمخرج من ASA، يجب أن تلبي الدالة التوقعات التالية:

  • يتوقع Azure Stream Analytics حالة HTTP 200 من تطبيق Functions للدفعات التي يعالجها بنجاح.
  • عندما يتلقى Azure Stream Analytics استثناء 413 ("http Request Entity Too Large") من دالة Azure، فإنه يقلل من حجم الدفعات التي يرسلها إلى دالة Azure.
  • أثناء الاتصال التجريبي، يرسل Stream Analytics طلب POST مع دفعة فارغة إلى دالات Azure ويتوقع أن يعود حالة HTTP 20x للتحقق من صحة الاختبار.

الخيار 1: التحديث بالمفتاح باستخدام ربط SQL للوظيفة Azure

يستخدم هذا الخيار ربط إخراج SQL الخاص بوظيفة Azure. يمكن لهذا الامتداد استبدال كائن في جدول دون الحاجة لكتابة عبارة SQL. في الوقت الحالي، لا يدعم عوامل التخصيص المركبة (التراكمات).

تم بناء هذه العينة على:

لفهم نهج الربط بشكل أفضل، تابع هذا الدرس.

أولاً، أنشئ تطبيق وظيفة HttpTrigger افتراضياً باتباع هذا البرنامج التعليمي. استخدم المعلومات التالية:

  • اللغة: C#‎
  • وقت التشغيل: .NET 6 (ضمن الوظيفة/وقت التشغيل v4)
  • القالب: HTTP trigger

قم بتثبيت ملحق الربط عن طريق تشغيل الأمر التالي في محطة طرفية موجودة في مجلد المشروع:

dotnet add package Microsoft.Azure.WebJobs.Extensions.Sql --prerelease

أضف عنصر SqlConnectionString في قسم Values في local.settings.json، وملء سلسلة الاتصال للخادم الوجهة:

{
    "IsEncrypted": false,
    "Values": {
        "AzureWebJobsStorage": "UseDevelopmentStorage=true",
        "FUNCTIONS_WORKER_RUNTIME": "dotnet",
        "SqlConnectionString": "Your connection string"
    }
}

استبدل الوظيفة بأكملها (ملف cs. في المشروع) بقصاصة التعليمة البرمجية التالية. حدث مساحة الاسم، اسم الفئة، واسم الوظيفة بنفس الاسم الخاص بك:

using System;
using System.IO;
using System.Threading.Tasks;
using Microsoft.AspNetCore.Mvc;
using Microsoft.Azure.WebJobs;
using Microsoft.Azure.WebJobs.Extensions.Http;
using Microsoft.AspNetCore.Http;
using Microsoft.Extensions.Logging;
using Newtonsoft.Json;

namespace Company.Function
{
    public static class HttpTrigger1{
        [FunctionName("HttpTrigger1")]
        public static async Task<IActionResult> Run (
            // http trigger binding
            [HttpTrigger(AuthorizationLevel.Function, "get","post", Route = null)] HttpRequest req,
            ILogger log,
            [Sql("dbo.device_updated", ConnectionStringSetting = "SqlConnectionString")] IAsyncCollector<Device> devices
            )
        {

            // Extract the body from the request
            string requestBody = await new StreamReader(req.Body).ReadToEndAsync();
            if (string.IsNullOrEmpty(requestBody)) {return new StatusCodeResult(204);} // 204, ASA connectivity check

            dynamic data = JsonConvert.DeserializeObject(requestBody);

            // Reject if too large, as per the doc
            if (data.ToString().Length > 262144) {return new StatusCodeResult(413);} //HttpStatusCode.RequestEntityTooLarge

            // Parse items and send to binding
            for (var i = 0; i < data.Count; i++)
            {
                var device = new Device();
                device.DeviceId = data[i].DeviceId;
                device.Value = data[i].Value;
                device.Timestamp = data[i].Timestamp;

                await devices.AddAsync(device);
            }
            await devices.FlushAsync();

            return new OkResult(); // 200
        }
    }

    public class Device{
        public int DeviceId { get; set; }
        public double Value { get; set; }
        public DateTime Timestamp { get; set; }
    }
}

قم بتحديث اسم الجدول الوجهة في قسم الربط:

[Sql("dbo.device_updated", ConnectionStringSetting = "SqlConnectionString")] IAsyncCollector<Device> devices

قم بتحديث فئة Device وقسم التعيين لمطابقة مخططك الخاص:

...
                device.DeviceId = data[i].DeviceId;
                device.Value = data[i].Value;
                device.Timestamp = data[i].Timestamp;
...
    public class Device{
        public int DeviceId { get; set; }
        public double Value { get; set; }
        public DateTime Timestamp { get; set; }

يمكنك الآن اختبار الأسلاك بين الدالة المحلية وقاعدة البيانات عن طريق تصحيح الأخطاء (F5 في تعليمة Visual Studio برمجية). يجب أن تكون قاعدة بيانات SQL قابلة للوصول من جهازك. يمكنك استخدام SSMS للتحقق من الاتصال. ثم أرسل طلبات POST إلى نقطة النهاية المحلية. يجب أن يعيد الطلب الذي يحتوي على متن فارغ HTTP 204. يجب الاستمرار في طلب مع حمولة فعلية في الجدول الوجهة (في وضع الاستبدال/التحديث). فيما يلي نموذج حمولة يتوافق مع المخطط المستخدم في هذا النموذج:

[{"DeviceId":3,"Value":13.4,"Timestamp":"2021-11-30T03:22:12.991Z"},{"DeviceId":4,"Value":41.4,"Timestamp":"2021-11-30T03:22:12.991Z"}]

يمكن الآن نشر الوظيفة في Azure. قم بتعيين إعداد تطبيق ل SqlConnectionString. يجب أن يسمح جدار الحماية Azure SQL Serverبخدمات Azure للوظيفة المباشرة للوصول إليها.

يمكنك بعد ذلك تعريف الدالة كمخرج في وظيفة ASA، واستخدامها لاستبدال السجلات بدلا من إدخالها.

الخيار 2: الدمج مع التعيين المركب (التراكم) عبر استعلام SQL مخصص

إشعار

عند إعادة التشغيل والاستعادة، قد يعيد ASA إرسال أحداث الإخراج التي أرسلها بالفعل. هذا السلوك يمكن أن يتسبب في فشل منطق التراكم (مضاعفة القيم الفردية). لمنع هذه المشكلة، قم بإخراج نفس البيانات في جدول باستخدام مخرج ASA SQL الأصلي. يمكنك استخدام جدول التحكم هذا لاكتشاف المشاكل وإعادة مزامنة التراكم عند الحاجة.

يستخدم هذا الخيار Microsoft.Data.SqlClient. تتيح لك هذه المكتبة إرسال أي استفسارات SQL إلى قاعدة بيانات SQL.

تم بناء هذه العينة على:

أولاً، أنشئ تطبيق وظيفة HttpTrigger افتراضياً باتباع هذا البرنامج التعليمي. يتم استخدام المعلومات التالية:

  • اللغة: C#‎
  • وقت التشغيل: .NET 6 (ضمن الوظيفة/وقت التشغيل v4)
  • القالب: HTTP trigger

قم بتثبيت مكتبة SqlClient عن طريق تشغيل الأمر التالي في محطة طرفية موجودة في مجلد المشروع:

dotnet add package Microsoft.Data.SqlClient --version 4.0.0

أضف عنصر SqlConnectionString في قسم Values في local.settings.json، وملء سلسلة الاتصال للخادم الوجهة:

{
    "IsEncrypted": false,
    "Values": {
        "AzureWebJobsStorage": "UseDevelopmentStorage=true",
        "FUNCTIONS_WORKER_RUNTIME": "dotnet",
        "SqlConnectionString": "Your connection string"
    }
}

استبدل الوظيفة بأكملها (ملف cs. في المشروع) بقصاصة التعليمة البرمجية التالية. قم بتحديث مساحة الاسم واسم الفئة واسم الوظيفة بنفسك:

using System;
using System.IO;
using System.Threading.Tasks;
using Microsoft.AspNetCore.Mvc;
using Microsoft.Azure.WebJobs;
using Microsoft.Azure.WebJobs.Extensions.Http;
using Microsoft.AspNetCore.Http;
using Microsoft.Extensions.Logging;
using Newtonsoft.Json;
using Microsoft.Data.SqlClient;

namespace Company.Function
{
    public static class HttpTrigger1{
        [FunctionName("HttpTrigger1")]
        public static async Task<IActionResult> Run(
            [HttpTrigger(AuthorizationLevel.Function, "get","post", Route = null)] HttpRequest req,
            ILogger log)
        {
            // Extract the body from the request
            string requestBody = await new StreamReader(req.Body).ReadToEndAsync();
            if (string.IsNullOrEmpty(requestBody)) {return new StatusCodeResult(204);} // 204, ASA connectivity check

            dynamic data = JsonConvert.DeserializeObject(requestBody);

            // Reject if too large, as per the doc
            if (data.ToString().Length > 262144) {return new StatusCodeResult(413);} //HttpStatusCode.RequestEntityTooLarge

            var SqlConnectionString = Environment.GetEnvironmentVariable("SqlConnectionString");
            using (SqlConnection conn = new SqlConnection(SqlConnectionString))
            {
                conn.Open();

                // Parse items and send to binding
                for (var i = 0; i < data.Count; i++)
                {
                    int DeviceId = data[i].DeviceId;
                    double Value = data[i].Value;
                    DateTime Timestamp = data[i].Timestamp;

                    var sqltext =
                    $"MERGE INTO [device_updated] AS old " +
                    $"USING (VALUES ({DeviceId},{Value},'{Timestamp}')) AS new (DeviceId, Value, Timestamp) " +
                    $"ON new.DeviceId = old.DeviceId " +
                    $"WHEN MATCHED THEN UPDATE SET old.Value += new.Value, old.Timestamp = new.Timestamp " +
                    $"WHEN NOT MATCHED BY TARGET THEN INSERT (DeviceId, Value, TimeStamp) VALUES (DeviceId, Value, Timestamp);";

                    //log.LogInformation($"Running {sqltext}");

                    using (SqlCommand cmd = new SqlCommand(sqltext, conn))
                    {
                        // Execute the command and log the # rows affected.
                        var rows = await cmd.ExecuteNonQueryAsync();
                        log.LogInformation($"{rows} rows updated");
                    }
                }
                conn.Close();
            }
            return new OkResult(); // 200
        }
    }
}

قم بتحديث قسم بناء الأوامر sqltext لمطابقة مخططك الخاص (لاحظ كيف يتم تحقيق التراكم عبر عامل التشغيل += عند التحديث):

    var sqltext =
    $"MERGE INTO [device_updated] AS old " +
    $"USING (VALUES ({DeviceId},{Value},'{Timestamp}')) AS new (DeviceId, Value, Timestamp) " +
    $"ON new.DeviceId = old.DeviceId " +
    $"WHEN MATCHED THEN UPDATE SET old.Value += new.Value, old.Timestamp = new.Timestamp " +
    $"WHEN NOT MATCHED BY TARGET THEN INSERT (DeviceId, Value, TimeStamp) VALUES (DeviceId, Value, Timestamp);";

يمكنك الآن اختبار الأسلاك بين الوظيفة المحلية وقاعدة البيانات عن طريق تصحيح الأخطاء (F5 في VS Code). يجب أن تكون قاعدة بيانات SQL قابلة للوصول من جهازك. يمكنك استخدام SSMS للتحقق من الاتصال. ثم أرسل طلبات POST إلى نقطة النهاية المحلية. يجب أن يعيد الطلب الذي يحتوي على متن فارغ HTTP 204. يجب الاستمرار في طلب حمولة فعلية في الجدول الوجهة (في وضع التجميع/الدمج). فيما يلي نموذج حمولة يتوافق مع المخطط المستخدم في هذا النموذج:

[{"DeviceId":3,"Value":13.4,"Timestamp":"2021-11-30T03:22:12.991Z"},{"DeviceId":4,"Value":41.4,"Timestamp":"2021-11-30T03:22:12.991Z"}]

يمكن الآن نشر الوظيفة في Azure. يجب تعيين إعداد التطبيق لـ SqlConnectionString. يجب أن يسمح جدار الحماية Azure SQL Serverبخدمات Azure للوظيفة المباشرة للوصول إليها.

يمكن بعد ذلك تعريف الوظيفة على أنها مخرجات في وظيفة ASA، واستخدامها لاستبدال السجلات بدلاً من إدراجها.

البدائل

خارج دالات Azure، يمكن لعدة طرق تحقيق النتيجة المتوقعة. يصف هذا القسم بعض هذه الطرق.

المعالجة اللاحقة في SQL Database الهدف

تعمل مهمة الخلفية بمجرد إدراج البيانات في قاعدة البيانات عبر مخرجات ASA SQL القياسية.

بالنسبة ل Azure SQL، استخدم INSTEAD OFمشغلات DML لاعتراض الأوامر INSERT التي يصدرها ASA.

CREATE TRIGGER tr_devices_updated_upsert ON device_updated INSTEAD OF INSERT
AS
BEGIN
	MERGE device_updated AS old
	
	-- In case of duplicates on the key below, use a subquery to make the key unique via aggregation or ranking functions
	USING inserted AS new
		ON new.DeviceId = old.DeviceId

	WHEN MATCHED THEN 
		UPDATE SET
			old.Value += new.Value, 
			old.Timestamp = new.Timestamp

	WHEN NOT MATCHED THEN
		INSERT (DeviceId, Value, Timestamp)
		VALUES (new.DeviceId, new.Value, new.Timestamp);  
END;

بالنسبة لـ Synapse SQL، يمكن لـ ASA الإدراج في جدول مرحلي. يمكن للمهمة المتكررة بعد ذلك تحويل البيانات حسب الحاجة إلى جدول وسيط. وأخيرا، يتم نقل البيانات إلى جدول الإنتاج.

المعالجة المسبقة في Azure Cosmos DB

Azure Cosmos DB يدعم UPSERT في الأصل. هنا، لا يمكن إلا الإضافة أو الاستبدال. يجب عليك إدارة التراكمات على جانب العميل في Azure Cosmos DB.

إذا تطابقت المتطلبات، يمكنك استبدال قاعدة بيانات SQL المستهدفة بمثيل Azure Cosmos DB. يتطلب هذا التغيير تغييرا مهما في بنية الحل بشكل عام.

بالنسبة ل Synapse SQL، يمكنك استخدام Azure Cosmos DB كطبقة وسيطة عبر Azure Synapse Link for Azure Cosmos DB. استخدم Azure Synapse Link لإنشاء مخزن تحليلي. يمكنك بعد ذلك الاستعلام عن مخزن البيانات هذا مباشرة في Synapse SQL.

مقارنة البدائل

كل نهج يقدم عروض قيمة وقدرات مختلفة:

نوع خيار الأوضاع قاعدة بيانات Azure SQL Azure Synapse Analytics
مرحلة ما بعد المعالجة
أزرار التشغيل استبدال، تراكم + غير متاح، لا تتوفر المشغلات في Synapse SQL
التدريج استبدال، تراكم + +
المعالجة المسبقة
دالات Azure استبدال، تراكم + - (أداء صف تلو الآخر)
استبدال Azure Cosmos DB الاستبدال ‏‫غير متوفر‬ ‏‫غير متوفر‬
Azure Cosmos DB Azure Synapse Link الاستبدال ‏‫غير متوفر‬ +

الحصول على الدعم

لمزيد من المساعدة، جرب صفحة أسئلة Microsoft Q&A الخاصة ب Azure Stream Analytics.

الخطوات التالية