Bemærk
Adgang til denne side kræver godkendelse. Du kan prøve at logge på eller ændre mapper.
Adgang til denne side kræver godkendelse. Du kan prøve at ændre mapper.
Fabric REST API'en leverer et serviceendepunkt til CRUD-operationer af Fabric-elementer. I denne tutorial gennemgår vi et end-to-end scenarie for, hvordan man opretter og opdaterer et Spark-jobdefinitionselement. Der er tre trin på højt niveau:
- Opret et Spark-jobdefinitionselement med en indledende tilstand.
- Upload hoveddefinitionsfilen og andre biblioteksfiler.
- Opdater Spark-jobdefinitionselementet med OneLake-URL'en til hoveddefinitionsfilen og andre bibliotekfiler.
Forudsætninger
- Et Microsoft Entra-token er nødvendigt for at få adgang til Fabric REST API'en. MSAL-biblioteket anbefales for at hente tokenet. Du kan få flere oplysninger under Understøttelse af godkendelsesflow i MSAL.
- Der kræves et lagertoken for at få adgang til OneLake-API'en. Du kan få flere oplysninger under MSAL til Python.
Opret et Spark-jobdefinitionselement med starttilstanden
Fabric REST API'en definerer et samlet endepunkt for CRUD-operationer af Fabric-elementer. Slutpunktet er https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items.
Vareoplysningerne er angivet i anmodningsdelen. Her er et eksempel på anmodningskroppen til oprettelse af et Spark-jobdefinitionselement:
{
"displayName": "SJDHelloWorld",
"type": "SparkJobDefinition",
"definition": {
"format": "SparkJobDefinitionV1",
"parts": [
{
"path": "SparkJobDefinitionV1.json",
"payload": "<REDACTED>",
"payloadType": "InlineBase64"
}
]
}
}
I dette eksempel hedder Spark-jobdefinitionselementet .SJDHelloWorld Feltet payload er det base64-kodede indhold i den detaljerede opsætning. Efter dekodning er indholdet:
{
"executableFile":null,
"defaultLakehouseArtifactId":"",
"mainClass":"",
"additionalLakehouseIds":[],
"retryPolicy":null,
"commandLineArguments":"",
"additionalLibraryUris":[],
"language":"",
"environmentArtifactId":null
}
Her er to hjælpefunktioner til at kode og afkode den detaljerede konfiguration:
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
Her er kodeudsnittet til at oprette et Spark-jobdefinitionselement:
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)
Upload hoveddefinitionsfilen og andre lib-filer
Der kræves et lagertoken for at uploade filen til OneLake. Her er en hjælpefunktion til at hente lagertokenet:
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']
Nu har vi oprettet et Spark-jobdefinitionselement. For at gøre det kørbart skal vi opsætte hoveddefinitionsfilen og de nødvendige egenskaber. Slutpunktet for overførsel af filen for dette SJD-element er https://onelake.dfs.fabric.microsoft.com/{workspaceId}/{sjditemid}. Den samme "workspaceId" som i forrige trin bør bruges. Værdien af "sjditemid" kunne findes i responskroppen for det forrige trin. Her er kodestykket til konfiguration af hoveddefinitionsfilen:
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())
Følg den samme proces for at uploade de andre lib-filer, hvis det er nødvendigt.
Opdater Spark-jobdefinitionselementet med OneLake-URL'en til hoveddefinitionsfilen og andre bibliotekfiler
Indtil nu har vi oprettet et Spark-jobdefinitionselement med en indledende tilstand og uploadet hoveddefinitionsfilen og andre bibliotekfiler. Det sidste trin er at opdatere Spark-jobdefinitionselementet for at sætte URL-egenskaberne for hoveddefinitionsfilen og andre bibliotekfiler. Endepunktet for opdatering af Spark-jobdefinitionselementet er https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid}. De samme "workspaceId" og "sjditemid" som tidligere trin bør bruges. Her er kodeuddraget til at opdatere Spark-jobdefinitionen:
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)
For at opsummere hele processen er både Fabric REST API og OneLake API nødvendige for at oprette og opdatere et Spark-jobdefinitionselement. Fabric REST API'en bruges til at oprette og opdatere Spark-jobdefinitionselementet. OneLake API'en bruges til at uploade hoveddefinitionsfilen og andre biblioteksfiler. Den primære definitionsfil og andre lib-filer uploades først til OneLake. Derefter sættes URL-egenskaberne for hoveddefinitionsfilen og andre bibliotekfiler i Spark-jobdefinitionselementet.