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 @@ -420,38 +420,28 @@ protected void createLogicalDatabases() {
protected String generateAndUploadConfig(String artifactName) {
LOG.info("Generating and uploading shard configuration...");
JSONObject config = new JSONObject();
config.put("configType", "dataflow");
JSONObject shardConfigBulk = new JSONObject();
JSONArray dataShards = new JSONArray();
JSONArray shardConfigs = new JSONArray();

int shardIdx = 0;
for (Map.Entry<String, List<String>> entry : requestedShardMap.entrySet()) {
String instanceName = entry.getKey();
String ip = instanceIpMap.get(instanceName);
List<String> dbNames = entry.getValue();

JSONObject dataShard = new JSONObject();
dataShard.put("dataShardId", instanceName);
dataShard.put("host", ip);
dataShard.put("port", port);
dataShard.put("user", username);
dataShard.put("password", password);

JSONArray databases = new JSONArray();
for (String dbName : dbNames) {
JSONObject db = new JSONObject();
db.put("dbName", dbName);
db.put("databaseId", String.format("%s%02d%s", "shard_", shardIdx, dbName));
db.put("refDataShardId", instanceName);
databases.put(db);
JSONObject shardConfig = new JSONObject();
shardConfig.put("logicalShardId", String.format("%s%02d_%s", "shard_", shardIdx, dbName));
Comment thread
pratickchokhani marked this conversation as resolved.
shardConfig.put("host", ip);
shardConfig.put("port", port);
shardConfig.put("user", username);
shardConfig.put("password", password);
shardConfig.put("dbName", dbName);
shardConfigs.put(shardConfig);
}
shardIdx++;
dataShard.put("databases", databases);
dataShards.put(dataShard);
}

shardConfigBulk.put("dataShards", dataShards);
config.put("shardConfigurationBulk", shardConfigBulk);
config.put("shardConfigs", shardConfigs);

String configContent = config.toString();
GcsArtifact artifact =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -161,8 +161,7 @@ public void setUp() throws IOException, InterruptedException {
"",
"",
databases);
createAndUploadBulkShardConfigToGcs(
new ArrayList<>(List.of(dataShard)), gcsResourceManager);
createAndUploadShardConfigToGcs(List.of(dataShard), gcsResourceManager);

// create pubsub manager
pubsubResourceManager = setUpPubSubResourceManager();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -236,55 +236,34 @@ protected void createAndUploadReverseMultiShardConfigToGcs(
gcsResourceManager.createArtifact("input/shard.json", shardFileContents);
}

protected void createAndUploadBulkShardConfigToGcs(
ArrayList<DataShard> dataShardsList, GcsResourceManager gcsResourceManager) {
JSONObject bulkConfig = new JSONObject();
bulkConfig.put("configType", "dataflow");

JSONObject shardConfigBulk = new JSONObject();

JSONObject schemaSourceJson = new JSONObject();
schemaSourceJson.put("dataShardId", "");
schemaSourceJson.put("host", "");
schemaSourceJson.put("user", "");
schemaSourceJson.put("password", "");
schemaSourceJson.put("port", "");
schemaSourceJson.put("dbName", "");
shardConfigBulk.put("schemaSource", schemaSourceJson);

JSONArray dataShardsArray = new JSONArray();
protected void createAndUploadShardConfigToGcs(
List<DataShard> dataShardsList, GcsResourceManager gcsResourceManager) {
JSONObject config = new JSONObject();
JSONArray shardConfigs = new JSONArray();

if (dataShardsList != null) {
for (DataShard shardData : dataShardsList) {
JSONObject shardJson = new JSONObject();

shardJson.put("dataShardId", shardData.dataShardId);
shardJson.put("logicalShardId", shardData.dataShardId);
shardJson.put("host", shardData.host);
shardJson.put("user", shardData.user);
shardJson.put("password", shardData.password);
shardJson.put("port", shardData.port);
shardJson.put("dbName", shardData.dbName);
shardJson.put("namespace", shardData.namespace);
shardJson.put("connectionProperties", shardData.connectionProperties);

JSONArray databasesArray = new JSONArray();

for (Database dbData : shardData.databases) {
JSONObject dbJson = new JSONObject();
dbJson.put("dbName", dbData.dbName);
dbJson.put("databaseId", dbData.databaseId);
dbJson.put("refDataShardId", dbData.refDataShardId);
databasesArray.put(dbJson);
if (shardData.namespace != null) {
shardJson.put("namespace", shardData.namespace);
}
if (shardData.connectionProperties != null) {
shardJson.put("connectionProperties", shardData.connectionProperties);
}
shardJson.put("databases", databasesArray);
dataShardsArray.put(shardJson);
shardConfigs.put(shardJson);
}
}
shardConfigBulk.put("dataShards", dataShardsArray);

bulkConfig.put("shardConfigurationBulk", shardConfigBulk);
String shardFileContents = bulkConfig.toString();
config.put("shardConfigs", shardConfigs);
String shardFileContents = config.toString();
LOG.info("Shard file contents: {}", shardFileContents);
gcsResourceManager.createArtifact("input/shard-bulk.json", shardFileContents);
gcsResourceManager.createArtifact("input/shard-config.json", shardFileContents);
}

protected void createAndUploadShardContextFileToGcs(
Expand Down Expand Up @@ -316,40 +295,44 @@ protected PipelineLauncher.LaunchInfo launchBulkDataflowJob(
Boolean multiSharded)
throws IOException {
// launch dataflow template
if (multiSharded) {
flexTemplateDataflowJobResourceManager =
FlexTemplateDataflowJobResourceManager.builder(jobName)
.withTemplateName("Sourcedb_to_Spanner_Flex")
.withTemplateModulePath("v2/sourcedb-to-spanner")
.addParameter("instanceId", spannerResourceManager.getInstanceId())
.addParameter("databaseId", spannerResourceManager.getDatabaseId())
.addParameter("projectId", PROJECT)
.addParameter("outputDirectory", "gs://" + artifactBucketName)
.addParameter("sessionFilePath", getGcsPath("input/session.json", gcsResourceManager))
.addParameter(
"sourceConfigURL", getGcsPath("input/shard-bulk.json", gcsResourceManager))
.addEnvironmentVariable(
"additionalExperiments", Collections.singletonList("disable_runner_v2"))
.build();
} else {
flexTemplateDataflowJobResourceManager =
FlexTemplateDataflowJobResourceManager.builder(jobName)
.withTemplateName("Sourcedb_to_Spanner_Flex")
.withTemplateModulePath("v2/sourcedb-to-spanner")
.addParameter("instanceId", spannerResourceManager.getInstanceId())
.addParameter("databaseId", spannerResourceManager.getDatabaseId())
.addParameter("projectId", PROJECT)
.addParameter("outputDirectory", "gs://" + artifactBucketName)
.addParameter("sessionFilePath", getGcsPath("input/session.json", gcsResourceManager))
.addParameter("sourceConfigURL", cloudSqlResourceManager.getUri())
.addParameter("username", cloudSqlResourceManager.getUsername())
.addParameter("password", cloudSqlResourceManager.getPassword())
.addParameter("jdbcDriverClassName", "com.mysql.jdbc.Driver")
.addEnvironmentVariable(
"additionalExperiments", Collections.singletonList("disable_runner_v2"))
.build();
if (!multiSharded) {
DataShard dataShard =
new DataShard(
"shard1",
cloudSqlResourceManager.getHost(),
cloudSqlResourceManager.getUsername(),
cloudSqlResourceManager.getPassword(),
String.valueOf(cloudSqlResourceManager.getPort()),
cloudSqlResourceManager.getDatabaseName(),
null,
"useSSL=false&allowPublicKeyRetrieval=true",
new ArrayList<>());

ArrayList<DataShard> shards = new ArrayList<>();
shards.add(dataShard);
createAndUploadShardConfigToGcs(shards, gcsResourceManager);
}

FlexTemplateDataflowJobResourceManager.Builder builder =
FlexTemplateDataflowJobResourceManager.builder(jobName)
.withTemplateName("Sourcedb_to_Spanner_Flex")
.withTemplateModulePath("v2/sourcedb-to-spanner")
.addParameter("instanceId", spannerResourceManager.getInstanceId())
.addParameter("databaseId", spannerResourceManager.getDatabaseId())
.addParameter("projectId", PROJECT)
.addParameter("outputDirectory", "gs://" + artifactBucketName)
.addParameter("sessionFilePath", getGcsPath("input/session.json", gcsResourceManager))
.addParameter(
"sourceConfigURL", getGcsPath("input/shard-config.json", gcsResourceManager))
.addEnvironmentVariable(
"additionalExperiments", Collections.singletonList("disable_runner_v2"));

if (!multiSharded) {
builder.addParameter("jdbcDriverClassName", "com.mysql.jdbc.Driver");

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

In the updated flow, the sharded and non-sharded flows are basically the same. Why do we need this special handling ?

}

flexTemplateDataflowJobResourceManager = builder.build();

// Run
PipelineLauncher.LaunchInfo jobInfo = flexTemplateDataflowJobResourceManager.launchJob();
assertThatPipeline(jobInfo).isRunning();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,54 +49,29 @@ public record DataShard(

public record Database(String dbName, String databaseId, String refDataShardId) {}

protected void createAndUploadBulkShardConfigToGcs(
protected void createAndUploadShardConfigToGcs(
List<DataShard> dataShardsList, GcsResourceManager gcsResourceManager) {
JSONObject bulkConfig = new JSONObject();
bulkConfig.put("configType", "dataflow");
JSONObject config = new JSONObject();
JSONArray shardConfigs = new JSONArray();

JSONObject shardConfigBulk = new JSONObject();

JSONObject schemaSourceJson = new JSONObject();
schemaSourceJson.put("dataShardId", "");
schemaSourceJson.put("host", "");
schemaSourceJson.put("user", "");
schemaSourceJson.put("password", "");
schemaSourceJson.put("port", "");
schemaSourceJson.put("dbName", "");
shardConfigBulk.put("schemaSource", schemaSourceJson);

JSONArray dataShardsArray = new JSONArray();
if (dataShardsList != null) {
for (DataShard shardData : dataShardsList) {
JSONObject shardJson = new JSONObject();

shardJson.put("dataShardId", shardData.dataShardId());
shardJson.put("logicalShardId", shardData.dataShardId());
shardJson.put("host", shardData.host());
shardJson.put("user", shardData.user());
shardJson.put("password", shardData.password());
shardJson.put("port", shardData.port());
shardJson.put("dbName", shardData.dbName());
shardJson.put("namespace", shardData.namespace());
shardJson.put("connectionProperties", shardData.connectionProperties());

JSONArray databasesArray = new JSONArray();

for (Database dbData : shardData.databases()) {
JSONObject dbJson = new JSONObject();
dbJson.put("dbName", dbData.dbName());
dbJson.put("databaseId", dbData.databaseId());
dbJson.put("refDataShardId", dbData.refDataShardId());
databasesArray.put(dbJson);
}
shardJson.put("databases", databasesArray);
dataShardsArray.put(shardJson);
shardConfigs.put(shardJson);
}
}
shardConfigBulk.put("dataShards", dataShardsArray);

bulkConfig.put("shardConfigurationBulk", shardConfigBulk);
String shardFileContents = bulkConfig.toString();
gcsResourceManager.createArtifact("input/shard-bulk.json", shardFileContents);
config.put("shardConfigs", shardConfigs);
String shardFileContents = config.toString();
gcsResourceManager.createArtifact("input/shard-config.json", shardFileContents);
}

protected PipelineLauncher.LaunchInfo launchBulkDataflowJob(
Expand Down Expand Up @@ -128,15 +103,25 @@ protected PipelineLauncher.LaunchInfo launchBulkDataflowJob(
builder.addParameter("sessionFilePath", getGcsPath("session.json", gcsResourceManager));
}

if (multiSharded) {
builder.addParameter(
"sourceConfigURL", getGcsPath("input/shard-bulk.json", gcsResourceManager));
} else {
builder.addParameter(
"sourceConfigURL",
cloudSqlResourceManager.getUri() + "?useSSL=false&allowPublicKeyRetrieval=true");
builder.addParameter("username", cloudSqlResourceManager.getUsername());
builder.addParameter("password", cloudSqlResourceManager.getPassword());
if (!multiSharded) {
DataShard dataShard =
new DataShard(
"shard1",
cloudSqlResourceManager.getHost(),
cloudSqlResourceManager.getUsername(),
cloudSqlResourceManager.getPassword(),
String.valueOf(cloudSqlResourceManager.getPort()),
cloudSqlResourceManager.getDatabaseName(),
null,
"useSSL=false&allowPublicKeyRetrieval=true",
Collections.emptyList());
createAndUploadShardConfigToGcs(Collections.singletonList(dataShard), gcsResourceManager);
}

builder.addParameter(
"sourceConfigURL", getGcsPath("input/shard-config.json", gcsResourceManager));

if (!multiSharded) {
builder.addParameter("jdbcDriverClassName", "com.mysql.cj.jdbc.Driver");
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -182,7 +182,7 @@ public void shardedMigrationAndValidationE2E() throws Exception {
mySQLResourceManager2.getDatabaseName(),
LOGICAL_SHARD_2,
LOGICAL_SHARD_2))));
createAndUploadBulkShardConfigToGcs(dataShards, gcsClient);
createAndUploadShardConfigToGcs(dataShards, gcsClient);

// 3. Launch Bulk Pipeline (SourceDbToSpanner) with multiSharded=true
String gcsOutputDirectory = "gs://" + artifactBucketName + "/" + testId;
Expand Down
6 changes: 2 additions & 4 deletions v2/sourcedb-to-spanner/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -65,9 +65,7 @@ mvn test
### Executing Template

#### Required Parameters
* **sourceConfigURL** (Configuration to connect to the source database): Can be the JDBC URL or the location of the sharding config. (Example: jdbc:mysql://10.10.10.10:3306/testdb or gs://test1/shard.conf). Refer to src/main/scripts/create_simple_shard_config.bash for steps to generate a shard configuration.
* **username** (username of the source database): The username which can be used to connect to the source database.
* **password** (username of the source database): The username which can be used to connect to the source database.
* **sourceConfigURL** (Source connection config file URL): The URL of the source connection config file. The file format is dependent on the source type. For Astra, it will point to an Astra connection config file ([sample](src/test/resources/SourceConfig/astra-connection-config.json)). For JDBC, it will point to a JDBC sharding config file ([sample](src/test/resources/SourceConfig/jdbc-shard-config.json)). For Cassandra, it will point to a Cassandra driver config file ([sample](src/test/resources/SourceConfig/cassandra-driver-config.conf)). This parameter is required. Refer to src/main/scripts/create_simple_shard_config.bash for steps to generate a shard configuration.
Comment thread
pratickchokhani marked this conversation as resolved.
* **instanceId** (Cloud Spanner Instance Id.): The destination Cloud Spanner instance.
* **databaseId** (Cloud Spanner Database Id.): The destination Cloud Spanner database.
* **projectId** (Cloud Spanner Project Id.): This is the name of the Cloud Spanner project.
Expand All @@ -89,7 +87,7 @@ export JOB_NAME="${IMAGE_NAME}-`date +%Y%m%d-%H%M%S-%N`"
gcloud dataflow flex-template run ${JOB_NAME} \
--project=${PROJECT} --region=us-central1 \
--template-file-gcs-location=${TEMPLATE_IMAGE_SPEC} \
--parameters sourceConfigURL="jdbc:mysql://<source_ip>:3306/<mysql_db_name>",username=<mysql user>,password=<mysql pass>,instanceId="<spanner instanceid>",databaseId="<spanner_database_id>",projectId="$PROJECT",outputDirectory=gs://<gcs-dir> \
--parameters sourceConfigURL="gs://<bucket-name>/source-config.json",instanceId="<spanner instanceid>",databaseId="<spanner_database_id>",projectId="$PROJECT",outputDirectory=gs://<gcs-dir> \
--additional-experiments=disable_runner_v2
```
#### Replaying DLQ entries.
Expand Down
Loading
Loading