【问题标题】:Azure Databricks: create audit trail for who ran what query at what momentAzure Databricks:为谁在什么时间运行什么查询创建审计跟踪
【发布时间】:2020-09-06 15:34:29
【问题描述】:

我们有一项审核要求,以便深入了解谁在 Azure Databricks 中的什么时间执行了什么查询。 Azure Databricks / Spark UI / Jobs 选项卡已经列出了执行的 Spark 作业,包括完成的查询和提交的时间。但它不包括执行查询的人。

  1. 是否有我们可以与 Azure Databricks 一起使用的 API 来查询 UI 中显示的这些 Spark 作业详细信息? (Databricks REST API 似乎没有提供这个,但也许我忽略了一些东西)
  2. 我们是否可以确定谁创建了 Spark 作业(使用 API)

谢谢, 格罗

【问题讨论】:

    标签: azure databricks azure-databricks


    【解决方案1】:

    1。访问 Spark API

    一个。驱动程序节点(内部)访问 Azure Databricks Spark api:

    import requests
    
    driverIp = spark.conf.get('spark.driver.host')
    port = spark.conf.get("spark.ui.port")
    url = F"http://{driverIp}:{port}/api/v1/applications"
    r = requests.get(url, timeout=3.0)
    r.status_code, r.text
    

    例如,如果您从公共 API 收到此错误消息: PERMISSION_DENIED: Traffic on this port is not permitted

    b.对 Azure Databricks Spark API 的外部访问:

    import requests
    import json
    """
      Program access to Databricks Spark UI. 
      
      Works external to Databricks environment or running within.
      Requires a Personal Access Token. Treat this like a password, do not store in a notebook. Please refer to the Secrets API.
      This Python code requires F string support.
    
    """
    
    # https://<databricks-host>/driver-proxy-api/o/0/<cluster_id>/<port>/api/v1/applications/<application-id-from-master-spark-ui>/stages/<stage-id>
    port = spark.conf.get("spark.ui.port")
    clusterId = spark.conf.get("spark.databricks.clusterUsageTags.clusterId")
    host = "eastus2.azuredatabricks.net"
    workspaceId = "999999999999111"  # follows the 'o=' in the databricks URLs or zero
    token = "dapideedeadbeefdeadbeefdeadbeef68ee3"  # Personal Access token
    
    url = F"https://{host}/driver-proxy-api/o/{workspaceId}/{clusterId}/{port}/api/v1/applications/?status=running"
    r = requests.get(url, auth=("token", token))
    
    # print Application list response
    print(r.status_code, r.text)
    
    applicationId = r.json()[0]['id'] # assumes only one response
    
    url = F"https://{host}/driver-proxy-api/o/{workspaceId}/{clusterId}/{port}/api/v1/applications/{applicationId}/jobs"
    r = requests.get(url, auth=("token", token))
    
    print(r.status_code, r.json())
    

    2。抱歉,不,暂时没有。

    集群日志将在您查看的位置,但用户身份不存在。

    投票和跟踪这个想法:https://ideas.databricks.com/ideas/DBE-I-313 如何进入创意门户:https://docs.databricks.com/ideas.html

    【讨论】:

    • 感谢道格拉斯的回复。相关问题,您知道是否可以访问 Spark Monitoring REST API(从/通过 Databricks)? (spark.apache.org/docs/latest/monitoring.html#rest-api)
    • 重新回答了问题的两个部分。
    • 谢谢。我尝试了这两个示例,访问内部 API 会导致 Connection refused 响应。访问外部 API 会导致 403 PERMISSION_DENIED: Traffic on this port is not permitted 是否有我们应该启用的配置设置以使访问成为可能?
    • 也许您的网络安全组受到限制,现在很多公司都这样做了。我不知道如何覆盖 403 错误,我在 Azure Databricks 上收到了同样的错误。策略是找到 spark driver java 进程并附加到相同的 IP 和端口。我在 tcp6 的 localhost 接口上尝试了方法 #1,但运气不佳。您可以使用%sh netstat -anltp 来寻找正确的接口、地址、协议和端口。
    【解决方案2】:

    根据Douglas 的回答,我想出了这个函数,可以在 Databricks Notebooks 中使用并获取有关缓存 RDD 的一些信息:

    我希望它有所帮助。我正在寻找有关此主题的 2h 的信息。

    def get_databricks_rdd_info():
    
        import requests, json
    
        # Get Spark Context
        sc = spark.sparkContext
        # Get App Id (Notebook is attached to it)
        app_id = sc._jsc.sc().applicationId()
        # Where is my driver
        driver_ip = spark.conf.get('spark.driver.host')
        port = spark.conf.get("spark.ui.port")
        # Compose the query to the Spark UI API
        url = f"http://{driver_ip}:{port}/api/v1/applications/{app_id}/storage/rdd"
    
        # Make request
        r = requests.get(url, timeout=3.0)
    
        if r.status_code == 200:
            # Compose results
            df = spark.createDataFrame([json.dumps(r) for r in r.json()], T.StringType())
            json_schema = spark.read.json(df.rdd.map(lambda row: row.value)).schema
            df = df.withColumn('value', F.from_json(F.col('value'), json_schema))
            df = df.selectExpr('value.*')
            
            # Generate summary
            df_summary = (df
                          .withColumn('Name', F.element_at(F.split(F.col('name'), ' '), -1))
                          .withColumn('Cached', F.round(F.lit(100) * F.col('numCachedPartitions')/F.col('numPartitions'), 2))
                          .withColumn('Memory GB', F.round(F.col('memoryUsed')*1e-9, 2))
                          .withColumn('Disk GB', F.round(F.col('diskUsed')*1e-9, 2))
                          .withColumnRenamed('numPartitions', '# Partitions')
    
                          .select([
                              'Name',
                              'id',
                              'Cached',
                              'Memory GB',
                              'Disk GB',
                              '# Partitions',
                          ]))
        else:
            print('Some error happened, code:', r.status_code)
            df = None
            df_summary = None
            
            
        return df, df_summary
    
    
    

    您可以在 Databricks Notebooks 中将其用作:

    df_rdd_info, df_summary = get_databricks_rdd_info()
    display(df_summary)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-06-08
      • 2011-08-25
      • 2020-01-01
      • 2016-09-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多