Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
Önemli
Olay kancaları için destek Genel Önizleme.
Olay kancalarını, olaylar bir işlem hattının olay günlüğüne kalıcı hale getirildiğinde çalıştırılacak özel Python geri çağırma işlevleri eklemek için kullanabilirsiniz. Özel izleme ve uyarı çözümleri uygulamak için olay kancalarını kullanabilirsiniz. Örneğin, belirli olaylar gerçekleştiğinde e-posta göndermek veya günlüğe yazmak veya işlem hattı olaylarını izlemek için üçüncü taraf çözümlerle tümleştirmek için olay kancalarını kullanabilirsiniz.
Tek bir bağımsız değişkeni kabul eden ve bu bağımsız değişkenin bir olayı temsil eden bir sözlük olduğu bir Python işleviyle bir etkinlik kancası tanımlayın. Ardından olay kancalarını bir işlem hattı için kaynak kodun parçası olarak ekleyin. İşlem hattında tanımlanan tüm olay kancaları, her işlem hattı güncelleştirmesi sırasında oluşturulan tüm olayları işlemeye çalışır. İşlem hattınız birden çok kaynak kodu dosyasından oluşuyorsa, tüm tanımlı olay kancaları işlem hattının tamamına uygulanır. Olay kancaları işlem hattınızın kaynak koduna dahil olsa da, işlem hattı grafiğine dahil değildir.
Etkinlik kancalarını, Hive meta veri deposuna veya Unity Kataloğu'na yayın yapan işlem hatlarıyla kullanabilirsiniz.
Uyarı
- Python, olay kancalarını tanımlamak için desteklenen tek dildir. SQL arabirimi kullanılarak uygulanan bir işlem hattındaki olayları işleyen özel Python işlevlerini tanımlamak için, özel işlevleri işlem hattının parçası olarak çalışan ayrı bir Python kaynak dosyasına ekleyin. python işlevleri, işlem hattı çalıştırıldığında işlem hattının tamamına uygulanır.
- Olay kancaları yalnızca maturity_level
STABLEolduğu olaylarda tetiklenir. - Olay kancaları, işlem hattı güncellemeleri ile zaman uyumsuz olarak yürütülürken, diğer olay kancaları ile zaman uyumlu olarak yürütülür. Bu, aynı anda yalnızca tek bir olay kancasının çalıştırıldığı ve diğer olay kancalarının, o anda çalışan olay kancası tamamlanana kadar beklediği anlamına gelir. Bir olay kancası süresiz olarak çalıştırılırsa, diğer tüm olay kancalarını engeller.
- Lakeflow işlem hatları, bir işlem hattı güncelleştirmesi sırasında yayımlanan her olay için her olay kancasını çalıştırmaya çalışır. Geciken olay kancalarının kuyruğa alınmış tüm olayları işlemek için yeterli zamana sahip olmasını sağlamak amacıyla işlem hattı, hesaplamasını sonlandırmadan önce yapılandırılamayan sabit bir süre bekler. Ancak, işlem sonlandırılmadan önce tüm olaylar sırasında tüm kancaların tetikleneceği garanti edilmez.
Etkinlik çengel işlemesini izleme
Bir güncellemenin olay kancalarının durumunu izlemek için işlem hattındaki hook_progress olay günlüğü türünü kullanın. Döngüsel bağımlılıkları önlemek için olay kancaları hook_progress olaylar için tetiklenmez.
Olay bağlantısı tanımlama
Olay kancası tanımlamak için on_event_hook dekoratörü kullanın:
@dp.on_event_hook(max_allowable_consecutive_failures=None)
def user_event_hook(event):
# Python code defining the event hook
max_allowable_consecutive_failures, bir olay kancasının devre dışı bırakılmadan önce art arda kaç kez başarısız olabileceğini açıklar. Bir olay kancası hatası, olay kancası bir istisna oluşturduğunda tanımlanır. Olay kancası devre dışı bırakılırsa, işlem hattı yeniden başlatılana kadar yeni olayları işlemez.
max_allowable_consecutive_failures
0 veya Nonedeğerinden büyük veya buna eşit bir tamsayı olmalıdır.
None değeri (varsayılan olarak atanır), olay kancası için izin verilen ardışık hata sayısı sınırı olmadığı ve olay kancasının hiçbir zaman devre dışı bırakılmadığı anlamına gelir.
Olay kancası hataları ve olay kancalarının devre dışı bırakılması olay günlüğünde hook_progress olaylar olarak izlenebilir.
Olay kancası işlevi, bu olay kancasını tetikleyen olayın sözlük gösterimi olan tam olarak bir parametreyi kabul eden bir Python işlevi olmalıdır. Olay kancası fonksiyonunun herhangi bir dönüş değeri görmezden gelinir.
Örnek: İşlenmek üzere belirli olayları seçme
Aşağıdaki örnekte, işlenmek üzere belirli olayları seçen bir olay kancası gösterilmektedir. Özellikle, bu örnek işlem hattı STOPPING olayları alınana kadar bekler ve ardından sürücü günlüklerine stdoutbir ileti kaydı yapar.
@dp.on_event_hook
def my_event_hook(event):
if (
event['event_type'] == 'update_progress' and
event['details']['update_progress']['state'] == 'STOPPING'
):
print('Received notification that update is stopping: ', event)
Örnek: Tüm olayları Slack kanalına gönderme
Aşağıdaki örnek, Slack API'sini kullanarak alınan tüm olayları Slack kanalına gönderen bir olay kancası uygular.
Bu örnekte, Slack API'sinde kimlik doğrulaması yapmak için gereken belirteci güvenli bir şekilde depolamak için Databricks gizli anahtar kullanılır.
from pyspark import pipelines as dp
import requests
# Get a Slack API token from a Databricks secret scope.
API_TOKEN = dbutils.secrets.get(scope="<secret-scope>", key="<token-key>")
@dp.on_event_hook
def write_events_to_slack(event):
res = requests.post(
url='https://slack.com/api/chat.postMessage',
headers={
'Content-Type': 'application/json',
'Authorization': 'Bearer ' + API_TOKEN,
},
json={
'channel': '<channel-id>',
'text': 'Received event:\n' + event,
}
)
Örnek: Ardışık dört hatadan sonra devre dışı bırakılacak şekilde bir olay kancası yapılandırın
Aşağıdaki örnek, ardışık olarak dört kez başarısız olursa devre dışı bırakılan bir olay kancasının nasıl yapılandırıldığını gösterir.
from pyspark import pipelines as dp
import random
def run_failing_operation():
raise Exception('Operation has failed')
# Allow up to 3 consecutive failures. After a 4th consecutive
# failure, this hook is disabled.
@dp.on_event_hook(max_allowable_consecutive_failures=3)
def non_reliable_event_hook(event):
run_failing_operation()
Örnek: Olay kancası olan işlem hattı
Aşağıdaki örnekte, bir işlem hattının kaynak koduna olay kancası ekleme gösterilmektedir. Bu, işlem hattıyla olay kancalarını kullanmanın basit ama eksiksiz bir örneğidir.
from pyspark import pipelines as dp
import requests
import json
import time
API_TOKEN = dbutils.secrets.get(scope="<secret-scope>", key="<token-key>")
SLACK_POST_MESSAGE_URL = 'https://slack.com/api/chat.postMessage'
DEV_CHANNEL = 'CHANNEL'
SLACK_HTTPS_HEADER_COMMON = {
'Content-Type': 'application/json',
'Authorization': 'Bearer ' + API_TOKEN
}
# Create a single dataset.
@dp.table
def test_dataset():
return spark.range(5)
# Definition of event hook to send events to a Slack channel.
@dp.on_event_hook
def write_events_to_slack(event):
res = requests.post(url=SLACK_POST_MESSAGE_URL, headers=SLACK_HTTPS_HEADER_COMMON, json = {
'channel': DEV_CHANNEL,
'text': 'Event hook triggered by event: ' + event['event_type'] + ' event.'
})