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.
Fabric REST API, Fabric öğelerinin CRUD işlemleri için bir hizmet uç noktası sağlar. Bu eğitimde, bir Spark iş tanımı öğesi oluşturma ve güncelleme senaryosunu uçtan uca anlatıyoruz. Üç üst düzey adım söz konusu olur:
- Başlangıç durumu olan bir Spark iş tanımı öğesi oluşturun.
- Ana tanım dosyasını ve diğer lib dosyalarını karşıya yükleyin.
- Spark iş tanımı öğesini ana tanım dosyasının OneLake URL'si ve diğer tür dosyalarla güncelleyin.
Önkoşullar
- Fabric REST API'sine erişmek için bir Microsoft Entra jetonu gereklidir. Belirteci almak için MSAL kitaplığı önerilir. Daha fazla bilgi için bkz . MSAL'de kimlik doğrulama akışı desteği.
- OneLake API'sine erişmek için bir depolama belirteci gereklidir. Daha fazla bilgi için bkz . Python için MSAL.
Başlangıç durumu ile bir Spark iş tanımı öğesi oluşturun
Fabric REST API, Fabric öğelerinin CRUD işlemleri için birleşik bir uç nokta tanımlar. Uç nokta şeklindedir https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items.
Öğe ayrıntıları istek gövdesinde belirtilir. İşte Spark iş tanımı öğesi oluşturmak için istek gövdesine bir örnek:
{
"displayName": "SJDHelloWorld",
"type": "SparkJobDefinition",
"definition": {
"format": "SparkJobDefinitionV1",
"parts": [
{
"path": "SparkJobDefinitionV1.json",
"payload": "<REDACTED>",
"payloadType": "InlineBase64"
}
]
}
}
Bu örnekte, Spark iş tanımı öğesi olarak adlandırılmıştır SJDHelloWorld. alanı payload , ayrıntılı kurulumun base64 kodlanmış içeriğidir. Kod çözme işlemi tamamlandıktan sonra içerik şu şekildedir:
{
"executableFile":null,
"defaultLakehouseArtifactId":"",
"mainClass":"",
"additionalLakehouseIds":[],
"retryPolicy":null,
"commandLineArguments":"",
"additionalLibraryUris":[],
"language":"",
"environmentArtifactId":null
}
Ayrıntılı kurulumu kodlamak ve kodunu çözmek için iki yardımcı işlev aşağıdadır:
import base64
def json_to_base64(json_data):
# Serialize the JSON data to a string
json_string = json.dumps(json_data)
# Encode the JSON string as bytes
json_bytes = json_string.encode('utf-8')
# Encode the bytes as Base64
base64_encoded = base64.b64encode(json_bytes).decode('utf-8')
return base64_encoded
def base64_to_json(base64_data):
# Decode the Base64-encoded string to bytes
base64_bytes = base64_data.encode('utf-8')
# Decode the bytes to a JSON string
json_string = base64.b64decode(base64_bytes).decode('utf-8')
# Deserialize the JSON string to a Python dictionary
json_data = json.loads(json_string)
return json_data
İşte Spark iş tanım öğesi oluşturmak için kod parçası:
import requests
bearerToken = "<REDACTED>" # Replace this token with the real AAD token
headers = {
"Authorization": f"Bearer {bearerToken}",
"Content-Type": "application/json" # Set the content type based on your request
}
payload = "<REDACTED>"
# Define the payload data for the POST request
payload_data = {
"displayName": "SJDHelloWorld",
"Type": "SparkJobDefinition",
"definition": {
"format": "SparkJobDefinitionV1",
"parts": [
{
"path": "SparkJobDefinitionV1.json",
"payload": payload,
"payloadType": "InlineBase64"
}
]
}
}
# Make the POST request with Bearer authentication
sjdCreateUrl = f"https://api.fabric.microsoft.com//v1/workspaces/{workspaceId}/items"
response = requests.post(sjdCreateUrl, json=payload_data, headers=headers)
Ana tanım dosyasını ve diğer lib dosyalarını karşıya yükleme
Dosyayı OneLake'e yüklemek için bir depolama belirteci gereklidir. Depolama belirtecini almak için bir yardımcı işlev aşağıdadır:
import msal
def getOnelakeStorageToken():
app = msal.PublicClientApplication(
"<REDACTED>", # This field should be the client ID
authority="https://login.microsoftonline.com/microsoft.com")
result = app.acquire_token_interactive(scopes=["https://storage.azure.com/.default"])
print(f"Successfully acquired AAD token with storage audience:{result['access_token']}")
return result['access_token']
Şimdi bir Spark iş tanımı öğesi oluşturuldu. Çalıştırılabilir hale getirmek için ana tanım dosyasını ve gerekli özellikleri ayarlamamız gerekir. Bu SJD öğesinin dosyasını karşıya yüklemek için uç nokta şudur: https://onelake.dfs.fabric.microsoft.com/{workspaceId}/{sjditemid}. Önceki adımdaki aynı "workspaceId" kullanılmalıdır. "sjditemid" değeri, önceki adımın yanıt gövdesinde bulunabilir. Ana tanım dosyasını ayarlamak için kod parçacığı aşağıdadır:
import requests
# Three steps are required: create file, append file, flush file
onelakeEndPoint = "https://onelake.dfs.fabric.microsoft.com/workspaceId/sjditemid" # Replace the ID of workspace and item with the right one
mainExecutableFile = "main.py" # The name of the main executable file
mainSubFolder = "Main" # The sub folder name of the main executable file. Don't change this value
onelakeRequestMainFileCreateUrl = f"{onelakeEndPoint}/{mainSubFolder}/{mainExecutableFile}?resource=file" # The URL for creating the main executable file via the 'file' resource type
onelakePutRequestHeaders = {
"Authorization": f"Bearer {onelakeStorageToken}", # The storage token can be achieved from the helper function above
}
onelakeCreateMainFileResponse = requests.put(onelakeRequestMainFileCreateUrl, headers=onelakePutRequestHeaders)
if onelakeCreateMainFileResponse.status_code == 201:
# Request was successful
print(f"Main File '{mainExecutableFile}' was successfully created in OneLake.")
# With the previous step, the main executable file is created in OneLake. Now we need to append the content of the main executable file
appendPosition = 0
appendAction = "append"
### Main File Append.
mainExecutableFileSizeInBytes = 83 # The size of the main executable file in bytes
onelakeRequestMainFileAppendUrl = f"{onelakeEndPoint}/{mainSubFolder}/{mainExecutableFile}?position={appendPosition}&action={appendAction}"
mainFileContents = "<REDACTED>" # The content of the main executable file, please replace this with the real content of the main executable file
mainExecutableFileSizeInBytes = 83 # The size of the main executable file in bytes, this value should match the size of the mainFileContents
onelakePatchRequestHeaders = {
"Authorization": f"Bearer {onelakeStorageToken}",
"Content-Type": "text/plain"
}
onelakeAppendMainFileResponse = requests.patch(onelakeRequestMainFileAppendUrl, data = mainFileContents, headers=onelakePatchRequestHeaders)
if onelakeAppendMainFileResponse.status_code == 202:
# Request was successful
print(f"Successfully accepted main file '{mainExecutableFile}' append data.")
# With the previous step, the content of the main executable file is appended to the file in OneLake. Now we need to flush the file
flushAction = "flush"
### Main File flush
onelakeRequestMainFileFlushUrl = f"{onelakeEndPoint}/{mainSubFolder}/{mainExecutableFile}?position={mainExecutableFileSizeInBytes}&action={flushAction}"
print(onelakeRequestMainFileFlushUrl)
onelakeFlushMainFileResponse = requests.patch(onelakeRequestMainFileFlushUrl, headers=onelakePatchRequestHeaders)
if onelakeFlushMainFileResponse.status_code == 200:
print(f"Successfully flushed main file '{mainExecutableFile}' contents.")
else:
print(onelakeFlushMainFileResponse.json())
Gerekirse diğer lib dosyalarını karşıya yüklemek için aynı işlemi izleyin.
Spark iş tanımı öğesini, ana tanım dosyasının OneLake URL'si ve diğer lib dosyalarıyla güncelleyin
Şimdiye kadar, bir başlangıç durumu olan bir Spark iş tanımı öğesi oluşturduk ve ana tanım dosyasını ile diğer tür dosyaları yükledik. Son adım, ana tanım dosyası ve diğer lib dosyalarının URL özelliklerini ayarlamak için Spark iş tanım öğesini güncellemektir. Spark iş tanımı öğesinin güncellenmesi için son nokta .https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid} Önceki adımlardaki aynı "workspaceId" ve "sjditemid" kullanılmalıdır. İşte Spark iş tanım öğesini güncellemek için kod parçası:
mainAbfssPath = f"abfss://{workspaceId}@onelake.dfs.fabric.microsoft.com/{sjditemid}/Main/{mainExecutableFile}" # The workspaceId and sjditemid are the same as previous steps, the mainExecutableFile is the name of the main executable file
libsAbfssPath = f"abfss://{workspaceId}@onelake.dfs.fabric.microsoft.com/{sjditemid}/Libs/{libsFile}" # The workspaceId and sjditemid are the same as previous steps, the libsFile is the name of the libs file
defaultLakehouseId = '<REDACTED>' # Replace this with the real default lakehouse ID
updateRequestBodyJson = {
"executableFile": mainAbfssPath,
"defaultLakehouseArtifactId": defaultLakehouseId,
"mainClass": "",
"additionalLakehouseIds": [],
"retryPolicy": None,
"commandLineArguments": "",
"additionalLibraryUris": [libsAbfssPath],
"language": "Python",
"environmentArtifactId": None}
# Encode the bytes as a Base64-encoded string
base64EncodedUpdateSJDPayload = json_to_base64(updateRequestBodyJson)
# Print the Base64-encoded string
print("Base64-encoded JSON payload for SJD Update:")
print(base64EncodedUpdateSJDPayload)
# Define the API URL
updateSjdUrl = f"https://api.fabric.microsoft.com//v1/workspaces/{workspaceId}/items/{sjditemid}/updateDefinition"
updatePayload = base64EncodedUpdateSJDPayload
payloadType = "InlineBase64"
path = "SparkJobDefinitionV1.json"
format = "SparkJobDefinitionV1"
Type = "SparkJobDefinition"
# Define the headers with Bearer authentication
bearerToken = "<REDACTED>" # Replace this token with the real AAD token
headers = {
"Authorization": f"Bearer {bearerToken}",
"Content-Type": "application/json" # Set the content type based on your request
}
# Define the payload data for the POST request
payload_data = {
"displayName": "sjdCreateTest11",
"Type": Type,
"definition": {
"format": format,
"parts": [
{
"path": path,
"payload": updatePayload,
"payloadType": payloadType
}
]
}
}
# Make the POST request with Bearer authentication
response = requests.post(updateSjdUrl, json=payload_data, headers=headers)
if response.status_code == 200:
print("Successfully updated SJD.")
else:
print(response.json())
print(response.status_code)
Tüm süreci özetlemek gerekirse, Spark iş tanımı öğesi oluşturmak ve güncellemek için hem Fabric REST API hem de OneLake API gereklidir. Fabric REST API, Spark iş tanım öğesini oluşturmak ve güncellemek için kullanılır. OneLake API'si, ana tanım dosyasını ve diğer lib dosyalarını karşıya yüklemek için kullanılır. Ana tanım dosyası ve diğer lib dosyaları önce OneLake'e yüklenir. Daha sonra ana tanım dosyasının ve diğer lib dosyalarının URL özellikleri Spark iş tanım öğesinde ayarlanır.