寻求自动化从SSAS Tabular提取数据至Databricks的Scala/Python方案
从SSAS Tabular自动化提取数据到Databricks DataFrame的方案
Python 实现方案
依赖准备
在Databricks集群上安装pyadomd库(Python连接SSAS的常用工具),可通过集群库管理界面安装,或在Notebook中执行:
%pip install pyadomd
核心代码示例
from pyadomd import Pyadomd import pandas as pd # 替换为你的SSAS连接信息 connection_string = "Provider=MSOLAP;Data Source=你的SSAS服务器地址;Catalog=你的Tabular数据库名称;" # 替换为目标数据的MDX查询语句 mdx_query = """ SELECT [Dim Date].[Date].[Date].Members ON ROWS, [Measures].[Sales Amount] ON COLUMNS FROM [你的Tabular模型名称] """ # 连接SSAS并提取数据 with Pyadomd(connection_string) as conn: with conn.cursor().execute(mdx_query) as cur: # 转换为Pandas DataFrame pd_df = pd.DataFrame(cur.fetchall(), columns=[col[0] for col in cur.description]) # 转为Spark DataFrame spark_df = spark.createDataFrame(pd_df) # 验证结果 spark_df.display()
关键注意事项:
- 确保Databricks集群与SSAS服务器网络连通,开放默认端口2383
- 若SSAS用Windows身份验证,连接字符串需追加
Integrated Security=SSPI;,且Databricks运行账号需具备SSAS访问权限
Scala 实现方案
依赖准备
在Databricks集群的库配置中添加olap4j依赖,Maven坐标为:org.olap4j:olap4j:1.2.0
核心代码示例
import org.olap4j.OlapConnection import org.olap4j.driver.xmla.XmlaOlap4jDriver import java.sql.DriverManager import org.apache.spark.sql.Row import org.apache.spark.sql.types.{StructType, StructField, StringType} // 替换为你的SSAS连接信息 val connectionString = "jdbc:xmla:Server=http://你的SSAS服务器地址/olap/msmdpump.dll;Catalog=你的Tabular数据库名称;" // 替换为目标数据的MDX查询语句 val mdxQuery = """ SELECT [Dim Date].[Date].[Date].Members ON ROWS, [Measures].[Sales Amount] ON COLUMNS FROM [你的Tabular模型名称] """ // 初始化连接并执行查询 Class.forName(classOf[XmlaOlap4jDriver].getName) val connection = DriverManager.getConnection(connectionString).asInstanceOf[OlapConnection] val statement = connection.createStatement() val resultSet = statement.executeOlapQuery(mdxQuery) // 构建Spark DataFrame的Schema和数据 val metaData = resultSet.getMetaData val columns = (0 until metaData.getColumnCount).map(i => StructField(metaData.getColumnName(i+1), StringType, nullable = true) ).toArray val schema = StructType(columns) val rows = collection.mutable.ArrayBuffer[Row]() while (resultSet.next()) { val rowValues = (0 until metaData.getColumnCount).map(i => resultSet.getString(i+1)).toArray rows.append(Row.fromSeq(rowValues)) } val sparkDF = spark.createDataFrame(spark.sparkContext.parallelize(rows), schema) // 展示结果 sparkDF.display() // 关闭资源 resultSet.close() statement.close() connection.close()
关键注意事项:
- 若SSAS启用Windows身份验证,连接字符串需追加
Integrated Security=SSPI;,且Databricks集群需配置Kerberos或AD认证传递身份 - 处理大型数据集时建议分批次提取,避免内存溢出
内容的提问来源于stack exchange,提问作者Mrityunjaya Naik
相关产品推荐
相关产品推荐

