【问题标题】:Load the ResultSet of query in dataframe using Spark / java使用 Spark / java 在数据框中加载查询的 ResultSet
【发布时间】:2020-10-30 09:46:42
【问题描述】:

我想在数据框 Spark 中加载选择查询的结果集。

我正在使用以下代码:

public static void func (Dataset <Row> df){
    df.repartition(20); //one connection per partition, see below

    df.foreachPartition((Iterator<Row> t) -> {
        Connection conn = DriverManager.getConnection("url",
                "root", "");

        conn.setAutoCommit(true);
        Statement statement = conn.createStatement();

        final int batchSize = 100000;
        int i = 0;
        while (t.hasNext()) {
            Row row = t.next();
            try {

              ResultSet query =   statement.executeQuery("SELECT * FROM zones WHERE zones.id IN ("
                        +"'"  + row.getAs("idZones")
                        + "'"+ ")  ");

     

            }  catch (SQLException e) {
                e.printStackTrace();
            } finally {
              
            }
        }

        statement.close();
        conn.close();


    });

}

有没有可能将 ResultSet 加载到数据框中?

我需要你的帮助

谢谢。

【问题讨论】:

    标签: java sql dataframe apache-spark


    【解决方案1】:

    如果我正确理解您的问题,您希望在数据框中加载 SQL 表。为此,您需要执行以下操作:

    1. 创建sparkSession的对象。
    2. 将您的 JDBC 连接放入 Properties 对象中。
    3. 通过read方法加载SQL表。
    4. 您可以根据加载的数据框应用相关过滤器。

    请在下面找到代码作为示例。

    import org.apache.spark.SparkConf;
    import org.apache.spark.sql.Dataset;
    import org.apache.spark.sql.Row;
    import org.apache.spark.sql.SparkSession;
    import org.apache.spark.sql.functions;
    
    import java.util.Properties;
    
    import lombok.extern.slf4j.Slf4j;
    
    @Slf4j
    public class ReadFromSQLTable {
        public static void main(String[] args) {
            String applicationName = ReadFromSQLTable.class.getName();
            SparkConf sparkConf = new SparkConf().setAppName(applicationName).setMaster("local[2]");
            // using Dataset<Row>
            SparkSession sparkSession = SparkSession
                    .builder()
                    .config(sparkConf)
                    .getOrCreate();
    
    
            Properties connectionProperties = new Properties();
    
            connectionProperties.put("user", "root"); // user name of your SQL database
            connectionProperties.put("password", "password"); // password of SQL
            connectionProperties.setProperty("driver", "com.mysql.cj.jdbc.Driver");
    // Name of the database that i am interacting with is `test`. You will find this as part of URL.
    // Name of table that I want to load is the `employee`
            Dataset<Row> employeeDetail = sparkSession.read().jdbc("jdbc:mysql://127.0.0.1:3306/test",
                    "employee", connectionProperties);
    
            log.error("Printing table detail");
            employeeDetail.show(); // to show the dataset loaded on the console
            long count = employeeDetail.count();
            System.out.println("The count is = " + count);
            Dataset<Row> employeeDetail2 = employeeDetail.filter("employee_number < 2");
            employeeDetail2.show();
    }
    

    您可以对这些数据框应用任何类型的操作,例如过滤器、选择或任何其他 SQL 操作。

    我正在本地系统中运行此代码。我在代码中添加了注释,以便您轻松理解。如果您有任何疑问,请告诉我。

    我希望这能让您对如何开始将 SQL 表加载为数据框有所了解。

    【讨论】:

    • 我就是想把sql查询的结果加载到dataframe中。
    • 您可以在加载的数据集上应用您想要的任何转换。您的 Select 查询与加载整个数据集然后在该数据集上应用 in 查询一样好。
    • 但是如果我想使用像 (ST_GeomFromText) 之类的函数,spark 上没有定义女巫,我必须使用 sql 查询。
    • Spark中有用户定义函数(UDF)的概念。你可以按照你想要的方式定义你自己的自定义函数,然后你就可以使用它们了。您可能会发现这很有用blog.cloudera.com/working-with-udfs-in-apache-spark
    猜你喜欢
    • 2020-11-01
    • 2017-08-19
    • 2018-03-07
    • 1970-01-01
    • 2020-06-07
    • 2017-08-19
    • 2020-07-13
    • 2016-03-07
    • 1970-01-01
    相关资源
    最近更新 更多