【问题标题】:Spark MySQL connector in javajava中的Spark MySQL连接器
【发布时间】:2020-08-08 19:42:06
【问题描述】:

我想连接 spark 和 MySQL。我尝试了以下代码:

public class Get_Data_From_MySQL implements Serializable {

private static final org.apache.log4j.Logger LOGGER = org.apache.log4j.Logger.getLogger(Get_Data_From_MySQL.class);

private static final String MYSQL_CONNECTION_URL = "jdbc:mysql://localhost:3306/test";
private static final String MYSQL_USERNAME = "root";
private static final String MYSQL_PWD = "";

private static final SparkSession sparkSession =
        SparkSession.builder().master("local[*]").appName("Spark2JdbcDs").getOrCreate();

public static void main(String[] args) {
    //JDBC connection properties
    final Properties connectionProperties = new Properties();
    connectionProperties.put("user", MYSQL_USERNAME);
    connectionProperties.put("password", MYSQL_PWD);
    connectionProperties.put("driver", "com.mysql.jdbc.Driver");
    final String dbTable =
            "(select age from employe";
      //Load MySQL query result as Dataset
    Dataset<Row> jdbcDF =
            sparkSession.read()
                    .jdbc(MYSQL_CONNECTION_URL, dbTable, "age", 10001, 499999, 10, connectionProperties);

我得到这个错误:

Exception in thread "main" java.sql.SQLException: Communication link failure: Mauvais 'handshake'
at com.mysql.jdbc.MysqlIO.doHandshake(MysqlIO.java:659)

在这一行:

Dataset<Row> jdbcDF =sparkSession.read().jdbc(MYSQL_CONNECTION_URL, dbTable, "age", 10001, 499999, 10, connectionProperties);

我检查了MySQL用户和密码,一切都正确。

谢谢。

【问题讨论】:

    标签: java mysql apache-spark intellij-idea


    【解决方案1】:

    您面临的问题可能意味着数据库根本无法访问。这个问题可能是由于任何原因导致的,可能是 JDBC URL 中的 IP 地址或主机名错误,或者本地 DNS 服务器无法识别 JDBC URL 中的主机名,或者端口号可能丢失,或者 DB 服务器已关闭。可能有很多原因,所以我建议您再次检查您的方法并检查所有内容。 通过终端连接并检查

    【讨论】:

    • 我通过更改 mysql connector 的版本解决了这个问题。谢谢
    • 实际上有一个语法错误,我通过一些修改能够解决它。数据库连接不是问题。再次感谢您的回复
    【解决方案2】:
      package JavaSpark.Javs.SQL;
    
       import java.util.Properties;
    
       import org.apache.spark.sql.Dataset;
       import org.apache.spark.sql.Row;
       import org.apache.spark.sql.SparkSession;
    
    
     public class sparkSqlMysql {
    
    
    private static final org.apache.log4j.Logger LOGGER = org.apache.log4j.Logger.getLogger(sparkSqlMysql.class);
    
    private static final SparkSession sparkSession =
            SparkSession.builder().master("local[*]").appName("Spark2JdbcDs").getOrCreate();    
    
    public static void main(String[] args) {
        //JDBC connection properties
        
        final Properties connectionProperties = new Properties();
        connectionProperties.put("user", "root");
        connectionProperties.put("password","mypassword");
        connectionProperties.put("driver", "com.mysql.jdbc.Driver");
         // Load MySQL query result as Dataset
        Dataset<Row> jdbcDF2 =
                sparkSession.read()
                        .jdbc("jdbc:mysql://localhost:3306/SQLprep", "customer", connectionProperties);
        
        jdbcDF2.show();
     }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-02-28
      • 2015-09-13
      • 2018-04-03
      • 1970-01-01
      • 2011-02-25
      • 2017-08-15
      • 1970-01-01
      相关资源
      最近更新 更多