Επεξεργασία συμβάντων με χρήση τελεστή SQL

Ένας τελεστής SQL, που ονομάζεται επίσης πρόγραμμα επεξεργασίας κώδικα SQL, είναι μια νέα δυνατότητα μετασχηματισμού δεδομένων στις ροές συμβάντων Microsoft Fabric. Οι τελεστές SQL παρέχουν μια εμπειρία επεξεργασίας κώδικα όπου μπορείτε εύκολα να ορίσετε τη δική σας προσαρμοσμένη λογική μετασχηματισμού δεδομένων χρησιμοποιώντας απλές παραστάσεις SQL. Αυτό το άρθρο περιγράφει τον τρόπο χρήσης ενός τελεστή SQL για μετασχηματισμούς δεδομένων σε μια ροή συμβάντων.

Note

Τα ονόματα τεχνουργημάτων ροής συμβάντων που περιλαμβάνουν χαρακτήρα υπογράμμισης (_) ή τελεία (.) δεν είναι συμβατά με τελεστές SQL. Για την καλύτερη δυνατή εμπειρία, δημιουργήστε μια νέα ροή συμβάντων χωρίς να χρησιμοποιήσετε χαρακτήρες υπογράμμισης ή τελείες στο όνομα του τεχνουργήματος.

Prerequisites

  • Πρόσβαση σε έναν χώρο εργασίας στη λειτουργία άδειας χρήσης εκχωρημένων πόρων Fabric ή στη δοκιμαστική λειτουργία άδειας χρήσης με δικαιώματα Συμβάλλοντα ή υψηλότερα.

Προσθήκη τελεστή SQL σε ροή συμβάντων

Για να εκτελέσετε λειτουργίες επεξεργασίας ροής στις ροές δεδομένων σας χρησιμοποιώντας έναν τελεστή SQL, προσθέστε έναν τελεστή SQL στη ροή συμβάντων χρησιμοποιώντας τις ακόλουθες οδηγίες:

  1. Δημιουργήστε μια νέα ροή συμβάντων. Στη συνέχεια, προσθέστε έναν τελεστή SQL σε αυτόν χρησιμοποιώντας μία από τις ακόλουθες επιλογές:

    • Στην κορδέλα, επιλέξτε Μετασχηματισμός συμβάντων και, στη συνέχεια, επιλέξτε SQL.

      Στιγμιότυπο οθόνης που εμφανίζει την επιλογή ενός τελεστή SQL στο μενού για τον μετασχηματισμό συμβάντων.

    • Στον καμβά, επιλέξτε Μετασχηματισμός συμβάντων ή προσθήκη προορισμού και, στη συνέχεια, επιλέξτε Κώδικας SQL.

      Στιγμιότυπο οθόνης που εμφανίζει την επιλογή ενός τελεστή SQL στη λίστα για τον μετασχηματισμό συμβάντων στον καμβά.

  2. Ένας νέος κόμβος SQL προστίθεται στη ροή συμβάντων σας. Επιλέξτε το εικονίδιο μολυβιού για να συνεχίσετε τη ρύθμιση του τελεστή SQL.

    Στιγμιότυπο οθόνης που εμφανίζει την επιλογή του εικονιδίου μολυβιού στον κόμβο τελεστή SQL.

  3. Στο τμήμα παραθύρου Κώδικας SQL , καθορίστε ένα μοναδικό όνομα για τον κόμβο τελεστή SQL στη ροή συμβάντων.

  4. Επεξεργαστείτε το ερώτημα στην περιοχή ερωτήματος ή επιλέξτε Επεξεργασία ερωτήματος για να εισέλθετε στην προβολή του προγράμματος επεξεργασίας κώδικα πλήρους οθόνης.

    Στιγμιότυπο οθόνης που εμφανίζει το πλαίσιο για την εισαγωγή ενός ονόματος λειτουργίας και το κουμπί για την επεξεργασία ενός ερωτήματος στο παράθυρο Κώδικας SQL.

  5. Η λειτουργία επεξεργασίας κώδικα πλήρους οθόνης διαθέτει ένα παράθυρο εξερεύνησης εισόδου/εξόδου στην αριστερή πλευρά. Η ενότητα επεξεργασίας κώδικα είναι ρυθμιζόμενη, ώστε να μπορείτε να αλλάξετε το μέγεθός της σύμφωνα με τις προτιμήσεις σας. Η ενότητα προεπισκόπησης στο κάτω μέρος σάς δίνει τη δυνατότητα να προβάλετε τόσο τα δεδομένα εισόδου όσο και το αποτέλεσμα δοκιμής του ερωτήματός σας.

    Στιγμιότυπο οθόνης που εμφανίζει το πλήρες πρόγραμμα επεξεργασίας SQL.

  6. Επιλέξτε το κείμενο στην ενότητα Έξοδοι και, στη συνέχεια, εισαγάγετε ένα όνομα για τον κόμβο προορισμού. Ο τελεστής SQL υποστηρίζει όλους τους προορισμούς Real-Time Intelligence, συμπεριλαμβανομένης μιας βάσης συμβάντων, μιας λίμνης, ενός ενεργοποιητή ή μιας ροής.

    Στιγμιότυπο οθόνης που εμφανίζει την περιοχή Έξοδοι με επιλεγμένο το κουμπί συν.

  7. Καθορίστε ένα ψευδώνυμο ή ένα όνομα για τον προορισμό εξόδου όπου γράφονται τα δεδομένα που υποβάλλονται σε επεξεργασία μέσω του τελεστή SQL.

    Στιγμιότυπο οθόνης που εμφανίζει το όνομα μιας εξόδου.

  8. Προσθέστε ερώτημα SQL για τον απαιτούμενο μετασχηματισμό δεδομένων.

    Μια ροή συμβάντων είναι χτισμένη πάνω από την Ανάλυση ροής του Azure και υποστηρίζει την ίδια σημασιολογία ερωτημάτων της γλώσσας ερωτημάτων της Ανάλυσης ροής. Για να μάθετε περισσότερα σχετικά με τη σύνταξη και τη χρήση, ανατρέξτε στα θέματα Ανάλυση ροής του Azure και Αναφορά γλώσσας ερωτήματος ροής συμβάντων.

    Ακολουθεί η βασική δομή ερωτημάτων:

    SELECT 
    
        column1, column2, ... 
    
    INTO 
    
        [output alias] 
    
    FROM 
    
        [input alias] 
    

    Αυτό το παράδειγμα ερωτήματος δείχνει την ανίχνευση υψηλών θερμοκρασιών σε ένα δωμάτιο κάθε λεπτό:

    
        SELECT 
        System.Timestamp AS WindowEnd, 
        roomId, 
        AVG(temperature) AS AvgTemp 
    INTO 
        output 
    FROM 
        input 
    GROUP BY 
        roomId, 
        TumblingWindow(minute, 1) 
    HAVING 
        AVG(temperature) > 75 
    

    Αυτό το παράδειγμα ερωτήματος δείχνει μια CASE δήλωση για την κατηγοριοποίηση της θερμοκρασίας:

    SELECT
        deviceId, 
        temperature, 
        CASE  
            WHEN temperature > 85 THEN 'High' 
            WHEN temperature BETWEEN 60 AND 85 THEN 'Normal' 
            ELSE 'Low' 
        END AS TempCategory 
    INTO 
        CategorizedTempOutput 
    FROM 
        SensorInput 
    
  9. Στην κορδέλα, χρησιμοποιήστε την εντολή Δοκιμή ερωτήματος για να επικυρώσετε τη λογική μετασχηματισμού. Τα αποτελέσματα του ερωτήματος δοκιμής εμφανίζονται στην καρτέλα Αποτέλεσμα δοκιμής .

    Στιγμιότυπο οθόνης που δείχνει ένα αποτέλεσμα δοκιμής.

  10. Όταν ολοκληρώσετε τη δοκιμή, επιλέξτε Αποθήκευση στην κορδέλα για να επιστρέψετε στον καμβά ροής συμβάντων.

    Στιγμιότυπο οθόνης που εμφανίζει την κορδέλα για ένα ερώτημα, συμπεριλαμβανομένων εντολών για τη δοκιμή του ερωτήματος και την αποθήκευση.

  11. Στο παράθυρο Κώδικας SQL , εάν το κουμπί Αποθήκευση είναι ενεργοποιημένο, επιλέξτε το για να αποθηκεύσετε τις ρυθμίσεις.

    Στιγμιότυπο οθόνης που εμφανίζει το παράθυρο Κώδικας SQL και το κουμπί Αποθήκευση.

  12. Ρυθμίστε τις παραμέτρους του προορισμού.

    Στιγμιότυπο οθόνης που εμφανίζει μια ολοκληρωμένη ροή συμβάντων.

Περισσότερα παραδείγματα

Τα παρακάτω παραδείγματα δείχνουν κοινά σενάρια ανάλυσης σε πραγματικό χρόνο που μπορείτε να εφαρμόσετε με τον τελεστή SQL.

Συνάθροιση πωλήσεων ανά λεπτό πόλης - Χρησιμοποιείται TumblingWindow για τον υπολογισμό σταθερών, μη επικαλυπτόμενων συνόλων πωλήσεων ενός λεπτού ομαδοποιημένων κατά πόλη:

SELECT
    System.Timestamp AS WindowEnd,
    city,
    SUM(salesAmount) AS TotalSales
INTO
    output
FROM
    input
GROUP BY
    city,
    TumblingWindow(minute, 1)

Ανίχνευση ριπής και bot - Χρησιμοποιείται HoppingWindow για τον εντοπισμό χρηστών που υποβάλλουν έναν ασυνήθιστα υψηλό αριθμό παραγγελιών μέσα σε ένα κυλιόμενο παράθυρο πέντε λεπτών, που αξιολογείται κάθε λεπτό:

SELECT
    System.Timestamp AS WindowEnd,
    userId,
    COUNT(*) AS OrderCount
INTO
    output
FROM
    input
GROUP BY
    userId,
    HoppingWindow(minute, 5, 1)
HAVING
    COUNT(*) > 10

Επισήμανση ανωμαλίας σε σχέση με μια κυλιόμενη γραμμή βάσης - Χρησιμοποιείται HoppingWindow για τον υπολογισμό ενός κυλιόμενου μέσου όρου και την επισήμανση συσκευών των οποίων η μέγιστη μετρική τιμή υπερβαίνει το διπλάσιο του μέσου όρου εντός του παραθύρου, γεγονός που υποδεικνύει μια πιθανή ανωμαλία:

SELECT
    System.Timestamp AS WindowEnd,
    deviceId,
    AVG(metricValue) AS RollingAvg,
    MAX(metricValue) AS CurrentMax
INTO
    output
FROM
    input
GROUP BY
    deviceId,
    HoppingWindow(minute, 10, 1)
HAVING
    MAX(metricValue) > 2 * AVG(metricValue)

Εγγραφή σε πολλούς προορισμούς από έναν χειριστή SQL

Με τον τελεστή SQL, μπορείτε να στείλετε δεδομένα σε πολλαπλές καταβόθρες εξόδου ή προορισμούς προσθέτοντας πολλούς INTO όρους στο ερώτημα SQL και ορίζοντας πολλαπλές εξόδους.

Ορισμός πολλαπλών εξόδων στο πρόγραμμα επεξεργασίας ερωτημάτων

  1. Επιλέξτε Επεξεργασία (εικονίδιο μολυβιού) στον κόμβο τελεστή SQL για να ανοίξετε το τμήμα παραθύρου Κώδικας SQL.

  2. Στο τμήμα παραθύρου Κώδικας SQL , επιλέξτε Επεξεργασία ερωτήματος για να ανοίξετε το πρόγραμμα επεξεργασίας κώδικα πλήρους οθόνης.

    Στιγμιότυπο οθόνης που εμφανίζει το παράθυρο Κώδικας SQL.

  3. Στο πρόγραμμα επεξεργασίας κώδικα πλήρους οθόνης, επιλέξτε + στην ενότητα Έξοδοι για να προσθέσετε μια νέα έξοδο. Επιλέξτε τον τύπο εξόδου της επιλογής σας. Δημιουργεί ένα ψευδώνυμο της εξόδου που μπορείτε να το χρησιμοποιήσετε σε ένα ερώτημα. Επιλέξτε το όνομα της εξόδου που δημιουργήθηκε και εισαγάγετε ένα όνομα της επιλογής σας.

    Στιγμιότυπο οθόνης που δείχνει το κουμπί για την προσθήκη εξόδου στον πλήρη επεξεργαστή SQL.

Χρησιμοποιήστε πολλαπλά SELECT ... Δηλώσεις INTO

Κάθε SELECT δήλωση μπορεί να γράψει σε διαφορετική έξοδο. Προσθέστε το ερώτημα για να γράψετε την έξοδο σε πολλούς προορισμούς.

Στο παρακάτω παράδειγμα ερωτήματος, η πρώτη SELECT πρόταση εγγράφει σε μια έξοδο με το όνομα RawArchive (τύπος: Lakehouse) και η δεύτερη SELECT πρόταση εγγράφει σε μια έξοδο με το όνομα AggregationResults (τύπος: Eventhouse).


-- Query 1: Archive all data to Lakehouse
SELECT *
INTO [RawArchive]
FROM [SQLDemoES-stream]

-- Query 2: Aggregate and filter data to create a real time dashboard to an Eventhouse
SELECT System.Timestamp() AS EventTime, COUNT(*) AS EventCount
INTO [AggregationResults]
FROM [SQLDemoES-stream]
GROUP BY TumblingWindow(minute, 1)
HAVING COUNT(*) > 100

Επαναχρησιμοποίηση ενδιάμεσης λογικής (βέλτιστη πρακτική)

Εάν θέλετε να αποφύγετε την επανάληψη της λογικής, χρησιμοποιήστε έναν όρο WITH και εξαπλωθείτε σε πολλαπλές εξόδους από εκεί. Στο παρακάτω παράδειγμα, η InputStream έκφραση κοινού πίνακα (CTE) ορίζεται για ανάγνωση από τη ροή εισόδου μία φορά και, στη συνέχεια, οι δύο SELECT δηλώσεις αναφέρονται στο InputStream CTE για εγγραφή σε διαφορετικές εξόδους. Αυτή η προσέγγιση είναι πιο αποτελεσματική επειδή αποφεύγει την ανάγνωση από τη ροή εισόδου πολλές φορές.

  1. Εισαγάγετε το ακόλουθο ερώτημα στο πρόγραμμα επεξεργασίας κώδικα SQL για ανάγνωση από τη ροή εισόδου μία φορά και εγγραφή σε πολλαπλές εξόδους.

    
    --Base query:  Reading input stream once
    With InputStream AS(
    SELECT * 
    FROM [SQLDemoES-stream] )
    
    -- Query 1: Archive all data to Lakehouse
    SELECT *
    INTO [RawArchive]
    FROM InputStream
    
    -- Query 2: Aggregate and filter data to create a real time dashboard to an Eventhouse
    SELECT System.Timestamp() AS EventTime, COUNT(*) AS EventCount
    INTO [AggregationResults]
    FROM InputStream
    GROUP BY TumblingWindow(minute, 1)
    HAVING COUNT(*) > 100
    
    
  2. Επιλέξτε Δοκιμή ερωτήματος για να επικυρώσετε το αποτέλεσμα του ερωτήματος. Κάθε έξοδος που ορίζεται στο ερώτημα έχει μια ξεχωριστή καρτέλα στον πίνακα Αποτελέσματα δοκιμής .

    Στιγμιότυπο οθόνης που δείχνει ένα παράδειγμα προσθήκης πολλαπλών ερωτημάτων προορισμού στο πλήρες πρόγραμμα επεξεργασίας SQL.

  3. Επιλέξτε Αποθήκευση για να αποθηκεύσετε το ερώτημα και να κλείσετε το πρόγραμμα επεξεργασίας.

    Στιγμιότυπο οθόνης που δείχνει το κουμπί Αποθήκευση στον πλήρη επεξεργαστή SQL.

  4. Επιλέξτε ξανά Αποθήκευση στο παράθυρο Πρόγραμμα επεξεργασίας SQL.

  5. Επιλέξτε κάθε κόμβο προορισμού που δημιουργήθηκε από τον τελεστή SQL και, στη συνέχεια, διαμορφώστε τις ρυθμίσεις προορισμού για κάθε έναν από αυτούς.

    Στιγμιότυπο οθόνης που εμφανίζει τους συνδέσμους διαμόρφωσης για κάθε κόμβο προορισμού.

  6. Αφού ολοκληρώσετε τη ρύθμιση παραμέτρων, η ροή συμβάντων σας θα πρέπει να μοιάζει με το παρακάτω παράδειγμα, όπου ο κόμβος-χειριστής SQL έχει δύο προορισμούς εξόδου.

    Στιγμιότυπο οθόνης που δείχνει ένα παράδειγμα τελεστή SQL με πολλαπλές εξόδους.

Ρύθμιση παραμέτρων πολιτικών ταξινόμησης συμβάντων στον τελεστή SQL

Με τον τελεστή SQL, μπορείτε να επεξεργαστείτε δεδομένα χρησιμοποιώντας χρόνο συμβάντος ή εφαρμογής. Από προεπιλογή, το Eventstream χρησιμοποιεί την ώρα άφιξης. Για να επεξεργαστείτε κατά χρόνο συμβάντος, πρέπει να το ρυθμίσετε ρητά χρησιμοποιώντας TIMESTAMP BY το ερώτημά σας.

Δείγμα εισόδου

{
    "deviceId": "device123",
    "temperature": 72,
    "eventTime": "2024-01-01T12:00:00Z"
}

Δείγμα ερωτήματος με χρήση χρόνου συμβάντος


SELECT
    deviceId,
    temperature,
    System.Timestamp() AS EventTimestamp
INTO
    Output
FROM
    Input
TIMESTAMP BY eventTime;

Μπορείτε επίσης να προσθέσετε όρια για συμβάντα καθυστερημένης άφιξης και εκτός παραγγελίας στις ρυθμίσεις για προχωρημένους του τελεστή SQL.

Στιγμιότυπο οθόνης που εμφανίζει τις ρυθμίσεις για προχωρημένους ενός τελεστή SQL.

Limitations

  • Ο τελεστής SQL έχει σχεδιαστεί για να συγκεντρώνει όλη τη λογική μετασχηματισμού σας. Ως αποτέλεσμα, δεν μπορείτε να το χρησιμοποιήσετε μαζί με άλλους ενσωματωμένους τελεστές στην ίδια διαδρομή επεξεργασίας. Δεν υποστηρίζεται επίσης σύνδεση πολλών τελεστών SQL σε μία μόνο διαδρομή.

  • Εάν προσθέσετε έναν τελεστή SQL στην τοπολογία σας, πρέπει να δημιουργήσετε νέους κόμβους προορισμού. Οι υπάρχοντες κόμβοι προορισμού δεν μπορούν να επαναχρησιμοποιηθούν με τον τελεστή SQL.