Skip to content

Commit 8bb114a

Browse files
committed
Replace RavenDB Hub/Sink replication with ETL Sink task
- Remove VerityReplicationHub and VerityReplicationSink (pull replication model) - Remove migration 005_ConfigureHubSink - Add VeritySinkEtlTask — push-based ETL via RavenDB ETL Sink - Add migration 005_ConfigureSinkEtl to register the ETL task - Minor whitespace cleanup in SecEdgarCompanyImporter
1 parent 8fdb554 commit 8bb114a

6 files changed

Lines changed: 97 additions & 109 deletions

File tree

src/RavenDB.Samples.Verity.Model/HubSinks/VerityReplicationHub.cs

Lines changed: 0 additions & 15 deletions
This file was deleted.

src/RavenDB.Samples.Verity.Model/HubSinks/VerityReplicationSink.cs

Lines changed: 0 additions & 16 deletions
This file was deleted.

src/RavenDB.Samples.Verity.Model/SecEdgarCompanyImporter.cs

Lines changed: 24 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -9,12 +9,12 @@ namespace RavenDB.Samples.Verity.Model;
99
public static class SecEdgarCompanyImporter
1010
{
1111
public record CompanyImportData(
12-
string PaddedCik,
13-
string Name,
14-
string? Sic,
15-
string? SicDescription,
12+
string PaddedCik,
13+
string Name,
14+
string? Sic,
15+
string? SicDescription,
1616
DateTime FiscalYearStart);
17-
17+
1818
private static readonly string[] FirstNames =
1919
[
2020
"James", "Mary", "Robert", "Patricia", "Michael",
@@ -37,10 +37,10 @@ public static async Task<CompanyImportData> FetchCompanyDataAsync(
3737
response.EnsureSuccessStatusCode();
3838

3939
await using var stream = await response.Content.ReadAsStreamAsync(ct);
40-
using var doc = await JsonDocument.ParseAsync(stream, cancellationToken: ct);
41-
var root = doc.RootElement;
40+
using var doc = await JsonDocument.ParseAsync(stream, cancellationToken: ct);
41+
var root = doc.RootElement;
4242

43-
var name = root.GetProperty("name").GetString() ?? paddedCik;
43+
var name = root.GetProperty("name").GetString() ?? paddedCik;
4444
var fiscalYearEnd = root.TryGetProperty("fiscalYearEnd", out var fye) ? fye.GetString() : null;
4545

4646
var fiscalYearStart = new DateTime(1, int.Parse(fiscalYearEnd!.Substring(0, 2)), 1, 0, 0, 0, DateTimeKind.Utc)
@@ -49,9 +49,9 @@ public static async Task<CompanyImportData> FetchCompanyDataAsync(
4949
fiscalYearStart = fiscalYearStart.AddYears(-1);
5050

5151
return new CompanyImportData(
52-
PaddedCik: paddedCik,
53-
Name: name,
54-
Sic: root.TryGetProperty("sic", out var sic) ? sic.GetString() : null,
52+
PaddedCik: paddedCik,
53+
Name: name,
54+
Sic: root.TryGetProperty("sic", out var sic) ? sic.GetString() : null,
5555
SicDescription: root.TryGetProperty("sicDescription", out var sicDesc) ? sicDesc.GetString() : null,
5656
FiscalYearStart: fiscalYearStart);
5757
}
@@ -66,43 +66,43 @@ public static async Task<Company> StoreCompanyAsync(
6666
CancellationToken ct = default)
6767
{
6868
var rng = Random.Shared;
69-
69+
7070
var companyId = Company.BuildId(data.Name);
7171

7272
if (await session.Advanced.ExistsAsync(companyId, ct))
7373
return (await session.LoadAsync<Company>(companyId, ct))!;
7474

7575
var company = new Company
7676
{
77-
Id = companyId,
78-
Name = data.Name,
79-
Cik = data.PaddedCik,
80-
Sic = data.Sic,
81-
SicDescription = data.SicDescription,
77+
Id = companyId,
78+
Name = data.Name,
79+
Cik = data.PaddedCik,
80+
Sic = data.Sic,
81+
SicDescription = data.SicDescription,
8282
FiscalYearStart = data.FiscalYearStart,
8383
};
8484

8585
await session.StoreAsync(company, companyId, ct);
8686

8787
var usedPairs = new HashSet<string>();
88-
var domain = data.Name.Replace(" ", "").Replace(",", "").Replace(".", "").ToLowerInvariant();
88+
var domain = data.Name.Replace(" ", "").Replace(",", "").Replace(".", "").ToLowerInvariant();
8989

9090
for (var i = 0; i < 2; i++)
9191
{
9292
string firstName, lastName;
9393
do
9494
{
9595
firstName = FirstNames[rng.Next(FirstNames.Length)];
96-
lastName = LastNames[rng.Next(LastNames.Length)];
96+
lastName = LastNames[rng.Next(LastNames.Length)];
9797
} while (!usedPairs.Add($"{firstName} {lastName}"));
9898

9999
await session.StoreAsync(new User
100100
{
101-
Id = User.BuildId(data.Name, firstName, lastName),
102-
CompanyId = companyId,
103-
Name = firstName,
104-
Surname = lastName,
105-
Email = $"{firstName.ToLower()}{lastName.ToLower()}@{domain}.com"
101+
Id = User.BuildId(data.Name, firstName, lastName),
102+
CompanyIds = [companyId],
103+
Name = firstName,
104+
Surname = lastName,
105+
Email = $"{firstName.ToLower()}{lastName.ToLower()}@{domain}.com"
106106
}, ct);
107107
}
108108

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,33 @@
1+
using Raven.Client.Documents.Operations.ETL;
2+
using RavenDB.Samples.Verity.Model; // żeby widzieć Company / Report
3+
4+
namespace RavenDB.Samples.Verity.Model.Tasks;
5+
6+
public static class VeritySinkEtlTask
7+
{
8+
public const string ConnectionStringName = "Verity Sink Connection";
9+
public const string TaskName = "VeritySinkEtlTask";
10+
11+
public static RavenEtlConfiguration Create() => new()
12+
{
13+
Name = TaskName,
14+
ConnectionStringName = ConnectionStringName,
15+
Disabled = false,
16+
Transforms =
17+
[
18+
new Transformation
19+
{
20+
Name = "Companies",
21+
Collections = [Company.Collection],
22+
Script = "loadToCompanies(this);"
23+
},
24+
new Transformation
25+
{
26+
Name = "Reports",
27+
Collections = [Report.Collection],
28+
Script = "if (this.Summary) loadToReports(this);"
29+
}
30+
]
31+
32+
};
33+
}

src/RavenDB.Samples.Verity.Setup/Migrations/005_ConfigureHubSink.cs

Lines changed: 0 additions & 54 deletions
This file was deleted.
Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,40 @@
1+
using Raven.Client.Documents.Operations.ConnectionStrings;
2+
using Raven.Client.Documents.Operations.ETL;
3+
using Raven.Client.ServerWide;
4+
using Raven.Client.ServerWide.Operations;
5+
using Raven.Migrations;
6+
using RavenDB.Samples.Verity.Model.Tasks;
7+
8+
namespace RavenDB.Samples.Verity.Setup.Migrations;
9+
10+
[Migration(5)]
11+
public sealed class ConfigureSinkEtl(MigrationContext context) : Migration
12+
{
13+
public override void Up()
14+
{
15+
try
16+
{
17+
DocumentStore.Maintenance.Server.Send(
18+
new CreateDatabaseOperation(new DatabaseRecord(Constants.DatabaseSinkName)));
19+
}
20+
catch (Exception ex) when (ex.Message.Contains("already exists")) { }
21+
22+
DocumentStore.Maintenance.Send(new PutConnectionStringOperation<RavenConnectionString>(new RavenConnectionString
23+
{
24+
Name = VeritySinkEtlTask.ConnectionStringName,
25+
Database = Constants.DatabaseSinkName,
26+
TopologyDiscoveryUrls = string.IsNullOrEmpty(context.HubServerInternalUrl)
27+
? DocumentStore.Urls
28+
: [context.HubServerInternalUrl]
29+
}));
30+
31+
// 3) ETL na źródle (Verity)
32+
DocumentStore.Maintenance.Send(new AddEtlOperation<RavenConnectionString>(VeritySinkEtlTask.Create()));
33+
}
34+
35+
public override void Down()
36+
{
37+
throw new NotSupportedException(
38+
"Rolling back the Sink ETL is not supported. Remove the ETL task and connection string via RavenDB Studio.");
39+
}
40+
}

0 commit comments

Comments
 (0)