Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -298,12 +298,12 @@ private static long getLengthOrPrecision(SourceColumnType sourceColumnType, long
* Ref: https://docs.oracle.com/en/database/oracle/oracle-database/19/sqlrf/Data-Types.html#GUID-D4EC7A0D-C119-4CB6-B6A2-EB0BCEDDBD35
* Size of BFILE can go upto 2gb, so setting to Interger.MAX_VALUE
*/
.put("BFILE", ResultSet::getBytes, valuePassThrough, Integer.MAX_VALUE)
.put("BFILE", ResultSet::getBytes, byteArrayToByteBuffer, Integer.MAX_VALUE)
/*
* Ref: https://docs.oracle.com/en/database/oracle/oracle-database/19/sqlrf/Data-Types.html#GUID-D4EC7A0D-C119-4CB6-B6A2-EB0BCEDDBD35
* Size of LONG RAW/LONG can go upto 2gb, so setting to Interger.MAX_VALUE
*/
.put("LONG RAW", ResultSet::getBytes, valuePassThrough, Integer.MAX_VALUE)
.put("LONG RAW", ResultSet::getBytes, byteArrayToByteBuffer, Integer.MAX_VALUE)
.put("LONG", ResultSet::getString, valuePassThrough, Integer.MAX_VALUE)
/*
* Ref: https://docs.oracle.com/en/database/oracle/oracle-database/19/sqlrf/Data-Types.html#GUID-D4EC7A0D-C119-4CB6-B6A2-EB0BCEDDBD35
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,10 @@
import java.nio.file.Path;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.ResultSet;
import java.sql.Statement;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
Expand Down Expand Up @@ -61,6 +63,7 @@
* environment setup and assertConditions.
*/
public class SourceDbToSpannerITBase extends JDBCBaseIT {
protected String testUsername = null;
private static final Logger LOG = LoggerFactory.getLogger(SourceDbToSpannerITBase.class);

public MySQLResourceManager setUpMySQLResourceManager() {
Expand All @@ -71,6 +74,80 @@ public CloudMySQLResourceManager setUpCloudMySQLResourceManager() {
return CloudMySQLResourceManager.builder(testName).build();
}

public org.apache.beam.it.jdbc.OracleResourceManager setUpOracleResourceManager() {
return org.apache.beam.it.jdbc.OracleResourceManager.builder(testName).build();
}

protected void loadOracleSQLFileResource(
JDBCResourceManager jdbcResourceManager, String resourcePath) throws Exception {
String sql =
String.join(
" ", Resources.readLines(Resources.getResource(resourcePath), StandardCharsets.UTF_8));
loadOracleSQLToJdbcResourceManager(jdbcResourceManager, sql);
}

protected void loadOracleSQLFileResource(
JDBCResourceManager jdbcResourceManager, String resourcePath, String targetUsername)
throws Exception {
String sql =
String.join(
" ", Resources.readLines(Resources.getResource(resourcePath), StandardCharsets.UTF_8));
loadOracleSQLToJdbcResourceManager(jdbcResourceManager, sql, targetUsername);
}

protected void loadOracleSQLToJdbcResourceManager(
JDBCResourceManager jdbcResourceManager, String sql) throws Exception {
loadOracleSQLToJdbcResourceManager(jdbcResourceManager, sql, "SYSTEM");
}

protected void loadOracleSQLToJdbcResourceManager(
JDBCResourceManager jdbcResourceManager, String sql, String targetUsername) throws Exception {
LOG.info("Loading Oracle sql to jdbc resource manager in schema {}", targetUsername);
try (Connection connection =
DriverManager.getConnection(
jdbcResourceManager.getUri(),
jdbcResourceManager.getUsername(),
jdbcResourceManager.getPassword())) {

// Ensure creation of tables occurs in the isolated namespace
if (!"SYSTEM".equalsIgnoreCase(targetUsername)) {
try (Statement stmt = connection.createStatement()) {
stmt.execute("ALTER SESSION SET CURRENT_SCHEMA = " + targetUsername);
}
}

// Preprocess SQL to handle multi-line statements and newlines
sql = sql.replaceAll("\r\n", " ").replaceAll("\n", " ");

// Split into individual statements based on -- SPLIT --
String[] statements = sql.split("-- SPLIT --");

// Execute each statement
try (Statement statement = connection.createStatement()) {
for (String stmt : statements) {
if (!stmt.trim().isEmpty()) {
LOG.info("Executing Oracle statement: {}", stmt);
statement.executeUpdate(stmt);
}
}
}
} catch (Exception e) {
LOG.info("failed to load SQL into database: {}", sql);
throw new Exception("Failed to load SQL into database", e);
}
LOG.info("Successfully loaded sql to jdbc resource manager");
}

protected String setupOracleIsolatedUser(JDBCResourceManager jdbcResourceManager) {
String testUsername =
"BULK_"
+ java.util.UUID.randomUUID().toString().replace("-", "").substring(0, 8).toUpperCase();
LOG.info("Creating isolated Oracle user: {}", testUsername);
jdbcResourceManager.runSQLUpdate("CREATE USER " + testUsername + " IDENTIFIED BY password");
jdbcResourceManager.runSQLUpdate("GRANT ALL PRIVILEGES TO " + testUsername);
return testUsername;
}

public PostgresResourceManager setUpPostgreSQLResourceManager() {
return PostgresResourceManager.builder(testName).build();
}
Expand Down Expand Up @@ -109,15 +186,39 @@ protected void loadSQLFileResource(JDBCResourceManager jdbcResourceManager, Stri
loadSQLToJdbcResourceManager(jdbcResourceManager, sql);
}

protected void loadSQLFileResource(
JDBCResourceManager jdbcResourceManager, String resourcePath, String targetUsername)
throws Exception {
String sql =
String.join(
" ", Resources.readLines(Resources.getResource(resourcePath), StandardCharsets.UTF_8));
loadSQLToJdbcResourceManager(jdbcResourceManager, sql, targetUsername);
}

protected void loadSQLToJdbcResourceManager(JDBCResourceManager jdbcResourceManager, String sql)
throws Exception {
LOG.info("Loading sql to jdbc resource manager with uri: {}", jdbcResourceManager.getUri());
try {
Connection connection =
DriverManager.getConnection(
jdbcResourceManager.getUri(),
jdbcResourceManager.getUsername(),
jdbcResourceManager.getPassword());
loadSQLToJdbcResourceManager(jdbcResourceManager, sql, "SYSTEM");
}

protected void loadSQLToJdbcResourceManager(
JDBCResourceManager jdbcResourceManager, String sql, String targetUsername) throws Exception {
LOG.info(
"Loading sql to jdbc resource manager with uri: {} as user: {}",
jdbcResourceManager.getUri(),
targetUsername);
try (Connection connection =
DriverManager.getConnection(
jdbcResourceManager.getUri(),
jdbcResourceManager.getUsername(),
jdbcResourceManager.getPassword())) {

// Ensure creation of tables occurs in the isolated namespace
if (!"SYSTEM".equalsIgnoreCase(targetUsername)
&& jdbcResourceManager instanceof org.apache.beam.it.jdbc.OracleResourceManager) {
try (Statement stmt = connection.createStatement()) {
stmt.execute("ALTER SESSION SET CURRENT_SCHEMA = " + targetUsername);
}
}

// Preprocess SQL to handle multi-line statements and newlines
sql = sql.replaceAll("\r\n", " ").replaceAll("\n", " ");
Expand All @@ -126,13 +227,14 @@ protected void loadSQLToJdbcResourceManager(JDBCResourceManager jdbcResourceMana
String[] statements = sql.split(";");

// Execute each statement
Statement statement = connection.createStatement();
for (String stmt : statements) {
if (!stmt.trim().isEmpty()) {
// Skip SELECT statements
if (!stmt.trim().toUpperCase().startsWith("SELECT")) {
LOG.info("Executing statement: {}", stmt);
statement.executeUpdate(stmt);
try (Statement statement = connection.createStatement()) {
for (String stmt : statements) {
if (!stmt.trim().isEmpty()) {
// Skip SELECT statements
if (!stmt.trim().toUpperCase().startsWith("SELECT")) {
LOG.info("Executing statement: {}", stmt);
statement.executeUpdate(stmt);
}
}
}
}
Expand Down Expand Up @@ -303,6 +405,7 @@ protected String createAndUploadShardConfigToGcs(
shard.setLogicalShardId("Shard1");
shard.setUser(jdbcResourceManager.getUsername());
shard.setPassword(jdbcResourceManager.getPassword());

if (jdbcResourceManager instanceof PostgresResourceManager pgRm) {
shard.setHost(pgRm.getHost());
shard.setPort(String.valueOf(pgRm.getPort()));
Expand All @@ -326,6 +429,14 @@ protected String createAndUploadShardConfigToGcs(
"Unsupported JDBC resource manager type: " + jdbcResourceManager.getClass().getName());
}

if (testUsername != null) {
shard.setNamespace(testUsername);
shard.setUser(testUsername);
shard.setPassword("password");
} else if (jdbcResourceManager instanceof org.apache.beam.it.jdbc.OracleResourceManager) {
shard.setNamespace(jdbcResourceManager.getUsername().toUpperCase());
}

if (jobParameters != null && jobParameters.containsKey("namespace")) {
shard.setNamespace(jobParameters.get("namespace"));
}
Expand Down Expand Up @@ -457,6 +568,9 @@ private String sqlDialectFrom(ResourceManager resourceManager) {
if (resourceManager instanceof PostgresResourceManager) {
return SQLDialect.POSTGRESQL.name();
}
if (resourceManager instanceof org.apache.beam.it.jdbc.OracleResourceManager) {
return "ORACLE";
}
return SQLDialect.MYSQL.name();
}

Expand All @@ -465,6 +579,9 @@ private String driverClassNameFrom(JDBCResourceManager jdbcResourceManager) {
if (jdbcResourceManager instanceof PostgresResourceManager) {
return Class.forName("org.postgresql.Driver").getCanonicalName();
}
if (jdbcResourceManager instanceof org.apache.beam.it.jdbc.OracleResourceManager) {
return "oracle.jdbc.OracleDriver";
}
return Class.forName("com.mysql.jdbc.Driver").getCanonicalName();
} catch (ClassNotFoundException e) {
throw new IllegalArgumentException(e);
Expand All @@ -479,4 +596,36 @@ protected PipelineOperator.Config.Builder wrapConfiguration(
}
return builder;
}

public static List<Map<String, Object>> runIsolatedSQLQuery(
org.apache.beam.it.jdbc.JDBCResourceManager jdbcResourceManager,
String testUsername,
String query) {
try (Connection connection =
DriverManager.getConnection(
jdbcResourceManager.getUri(),
jdbcResourceManager.getUsername(),
jdbcResourceManager.getPassword());
Statement stmt = connection.createStatement()) {
if (!"SYSTEM".equalsIgnoreCase(testUsername)
&& jdbcResourceManager instanceof org.apache.beam.it.jdbc.OracleResourceManager) {
stmt.execute("ALTER SESSION SET CURRENT_SCHEMA = " + testUsername);
}
List<Map<String, Object>> result = new ArrayList<>();
try (ResultSet rs = stmt.executeQuery(query)) {
java.sql.ResultSetMetaData md = rs.getMetaData();
int columns = md.getColumnCount();
while (rs.next()) {
Map<String, Object> row = new HashMap<>(columns);
for (int i = 1; i <= columns; ++i) {
row.put(md.getColumnName(i).toLowerCase(), rs.getObject(i));
}
result.add(row);
}
}
return result;
} catch (Exception e) {
throw new RuntimeException("Error running isolated query", e);
}
}
}
Loading
Loading