Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
L'API Fabric REST fornisce un endpoint di servizio per le operazioni CRUD degli elementi Fabric. In questo tutorial, presentiamo uno scenario end-to-end su come creare e aggiornare un elemento di definizione del lavoro di Spark. Sono coinvolti tre passaggi di alto livello:
- Crea un elemento di definizione del lavoro Spark con uno stato iniziale.
- Caricare il file di definizione principale e altri file lib.
- Aggiorna l'elemento di definizione del lavoro Spark con l'URL OneLake del file principale di definizione e altri file di liberazione.
Prerequisiti
- Per accedere all'API REST di Fabric, è necessario un token Microsoft Entra. La libreria MSAL è consigliata per ottenere il token. Per altre informazioni, vedere Supporto del flusso di autenticazione in MSAL.
- Per accedere all'API OneLake, è necessario un token di archiviazione. Per altre informazioni, vedere MSAL per Python.
Crea un elemento di definizione del lavoro Spark con lo stato iniziale
L'API Fabric REST definisce un endpoint unificato per le operazioni CRUD degli elementi Fabric. L'endpoint è https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items.
I dettagli dell'elemento vengono specificati all'interno del corpo della richiesta. Ecco un esempio del corpo della richiesta per creare un elemento di definizione del lavoro Spark:
{
"displayName": "SJDHelloWorld",
"type": "SparkJobDefinition",
"definition": {
"format": "SparkJobDefinitionV1",
"parts": [
{
"path": "SparkJobDefinitionV1.json",
"payload": "<REDACTED>",
"payloadType": "InlineBase64"
}
]
}
}
In questo esempio, l'elemento di definizione del lavoro Spark è chiamato SJDHelloWorld. Il payload campo è il contenuto codificato in base64 della configurazione dettagliata. Dopo la decodifica, il contenuto è:
{
"executableFile":null,
"defaultLakehouseArtifactId":"",
"mainClass":"",
"additionalLakehouseIds":[],
"retryPolicy":null,
"commandLineArguments":"",
"additionalLibraryUris":[],
"language":"",
"environmentArtifactId":null
}
Di seguito sono indicate due funzioni helper per codificare e decodificare la configurazione dettagliata:
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
Ecco il snippet di codice per creare un elemento di definizione del lavoro Spark:
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)
Caricare il file di definizione principale e altri file lib
Per caricare il file in OneLake, è necessario un token di archiviazione. Di seguito è indicata una funzione helper per ottenere il token di archiviazione:
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']
Ora abbiamo creato un elemento di definizione del lavoro Spark. Per renderlo eseguibile, è necessario configurare il file di definizione principale e le proprietà necessarie. L'endpoint per il caricamento del file per questo elemento SJD è https://onelake.dfs.fabric.microsoft.com/{workspaceId}/{sjditemid}. È necessario usare lo stesso "workspaceId" del passaggio precedente. Il valore di "sjditemid" si poteva trovare nel corpo della risposta del passaggio precedente. Di seguito è riportato un frammento di codice per configurare il file di definizione principale:
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())
Seguire lo stesso processo per caricare gli altri file lib, se necessario.
Aggiorna l'elemento di definizione del lavoro Spark con l'URL OneLake del file principale di definizione e altri file di lib
Fino ad ora, abbiamo creato un elemento di definizione del lavoro Spark con uno stato iniziale e caricato il file principale di definizione e altri file di liberazione. L'ultimo passo è aggiornare l'elemento di definizione del lavoro di Spark per impostare le proprietà URL del file principale di definizione e degli altri file di liberazione. L'endpoint per aggiornare l'elemento di definizione del lavoro Spark è https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid}. Dovrebbero essere usati gli stessi "workspaceId" e "sjditemid" dei passaggi precedenti. Ecco il snippet di codice per aggiornare l'elemento di definizione del lavoro Spark:
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)
Per riassumere l'intero processo, sono necessarie sia l'API Fabric REST che l'API OneLake per creare e aggiornare un elemento di definizione del lavoro Spark. L'API Fabric REST viene utilizzata per creare e aggiornare l'elemento di definizione del lavoro Spark. L'API OneLake viene usata per caricare il file di definizione principale e altri file lib. Il file di definizione principale e altri file lib vengono prima caricati in OneLake. Poi le proprietà URL del file di definizione principale e degli altri file di lib vengono impostate nell'elemento di definizione del lavoro di Spark.