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
35 changes: 21 additions & 14 deletions data-pipeline/src/main.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import "dotenv/config";
import { task } from "@renderinc/sdk/workflows";
import { task, type TaskContext } from "@renderinc/sdk/workflows";

interface User {
id: string;
Expand Down Expand Up @@ -58,7 +58,7 @@ function simpleHash(str: string): number {

const fetchUserData = task(
{ name: "fetchUserData", retry },
async function fetchUserData(userIds: string[]) {
async function fetchUserData(_ctx: TaskContext, userIds: string[]) {
console.log(`[SOURCE] Fetching user data for ${userIds.length} users`);

const mockUsers: { [key: string]: User } = {
Expand All @@ -79,7 +79,7 @@ const fetchUserData = task(

const fetchTransactionData = task(
{ name: "fetchTransactionData", retry },
async function fetchTransactionData(userIds: string[], days: number = 30) {
async function fetchTransactionData(_ctx: TaskContext, userIds: string[], days: number = 30) {
console.log(`[SOURCE] Fetching transactions for ${userIds.length} users (${days} days)`);

const transactions: Transaction[] = [];
Expand Down Expand Up @@ -108,7 +108,7 @@ const fetchTransactionData = task(

const fetchEngagementData = task(
{ name: "fetchEngagementData", retry },
async function fetchEngagementData(userIds: string[]) {
async function fetchEngagementData(_ctx: TaskContext, userIds: string[]) {
console.log(`[SOURCE] Fetching engagement data for ${userIds.length} users`);

const engagement: Engagement[] = userIds.map((userId) => {
Expand All @@ -135,7 +135,7 @@ const fetchEngagementData = task(

const enrichWithGeoData = task(
{ name: "enrichWithGeoData", retry },
async function enrichWithGeoData(userEmail: string) {
async function enrichWithGeoData(_ctx: TaskContext, userEmail: string) {
console.log(`[ENRICH] Enriching geo data for ${userEmail}`);
const idx = simpleHash(userEmail) % 4;
return {
Expand All @@ -149,6 +149,7 @@ const enrichWithGeoData = task(
const calculateUserMetrics = task(
{ name: "calculateUserMetrics", retry },
async function calculateUserMetrics(
_ctx: TaskContext,
user: User,
transactions: Transaction[],
engagement: Engagement,
Expand Down Expand Up @@ -197,6 +198,7 @@ const calculateUserMetrics = task(
const transformUserData = task(
{ name: "transformUserData", retry },
async function transformUserData(
ctx: TaskContext,
userData: { data: User[] },
transactionData: { data: Transaction[] },
engagementData: { data: Engagement[] },
Expand All @@ -217,8 +219,8 @@ const transformUserData = task(
users.map(async (user) => {
const userEngagement = engagementMap.get(user.id) ?? ({} as Engagement);
const [userMetrics, geoData] = await Promise.all([
calculateUserMetrics(user, transactions, userEngagement),
enrichWithGeoData(user.email),
ctx.step(calculateUserMetrics, user, transactions, userEngagement),
ctx.step(enrichWithGeoData, user.email),
]);
return { ...userMetrics, geo: geoData };
}),
Expand All @@ -233,7 +235,7 @@ const transformUserData = task(

const aggregateInsights = task(
{ name: "aggregateInsights", retry },
function aggregateInsights(enrichedData: { data: EnrichedUser[] }) {
function aggregateInsights(_ctx: TaskContext, enrichedData: { data: EnrichedUser[] }) {
console.log("[AGGREGATE] Generating insights from enriched data");

const users = enrichedData.data ?? [];
Expand Down Expand Up @@ -284,7 +286,7 @@ const aggregateInsights = task(
// Root task: full pipeline orchestrator
task(
{ name: "runDataPipeline", retry, timeoutSeconds: 300 },
async function runDataPipeline(userIds: string[]) {
async function runDataPipeline(ctx: TaskContext, userIds: string[]) {
console.log("=".repeat(80));
console.log("[PIPELINE] Starting Data Pipeline");
console.log(`[PIPELINE] Processing ${userIds.length} users`);
Expand All @@ -293,9 +295,9 @@ task(
// Stage 1: Parallel extraction
console.log("[PIPELINE] Stage 1/3: EXTRACT (parallel)");
const [userData, transactionData, engagementData] = await Promise.all([
fetchUserData(userIds),
fetchTransactionData(userIds),
fetchEngagementData(userIds),
ctx.step(fetchUserData, userIds),
ctx.step(fetchTransactionData, userIds),
ctx.step(fetchEngagementData, userIds),
]);

console.log(
Expand All @@ -304,12 +306,17 @@ task(

// Stage 2: Transform
console.log("[PIPELINE] Stage 2/3: TRANSFORM");
const enrichedData = await transformUserData(userData, transactionData, engagementData);
const enrichedData = await ctx.step(
transformUserData,
userData,
transactionData,
engagementData,
);
console.log(`[PIPELINE] Enriched ${enrichedData.count} user profiles`);

// Stage 3: Aggregate
console.log("[PIPELINE] Stage 3/3: AGGREGATE");
const insights = await aggregateInsights(enrichedData as { data: EnrichedUser[] });
const insights = await ctx.step(aggregateInsights, enrichedData);

const pipelineResult = {
status: "success",
Expand Down
20 changes: 10 additions & 10 deletions etl-job/src/main.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import "dotenv/config";
import { task } from "@renderinc/sdk/workflows";
import { task, type TaskContext } from "@renderinc/sdk/workflows";
import { readFileSync, existsSync } from "node:fs";
import { resolve } from "node:path";

Expand Down Expand Up @@ -32,7 +32,7 @@ const retry = {
// Subtask: extract rows from a CSV file
const extractCsvData = task(
{ name: "extractCsvData", retry },
function extractCsvData(filePath: string): Record[] {
function extractCsvData(_ctx: TaskContext, filePath: string): Record[] {
console.log(`[EXTRACT] Reading CSV file: ${filePath}`);

const fullPath = resolve(filePath);
Expand Down Expand Up @@ -69,7 +69,7 @@ const extractCsvData = task(
// Subtask: validate and clean a single record
const validateRecord = task(
{ name: "validateRecord", retry },
function validateRecord(record: Record): ValidatedRecord {
function validateRecord(_ctx: TaskContext, record: Record): ValidatedRecord {
console.log(`[TRANSFORM] Validating record ID: ${record.id ?? "unknown"}`);

const errors: string[] = [];
Expand Down Expand Up @@ -112,15 +112,15 @@ const validateRecord = task(
// Subtask: validate a batch of records by calling validateRecord for each
const transformBatch = task(
{ name: "transformBatch", retry },
async function transformBatch(records: Record[]) {
async function transformBatch(ctx: TaskContext, records: Record[]) {
console.log(`[TRANSFORM] Starting batch transformation of ${records.length} records`);

const validRecords: ValidatedRecord[] = [];
const invalidRecords: ValidatedRecord[] = [];

for (let i = 0; i < records.length; i++) {
console.log(`[TRANSFORM] Processing record ${i + 1}/${records.length}`);
const validated = await validateRecord(records[i]);
const validated = await ctx.step(validateRecord, records[i]);

if (validated.is_valid) {
validRecords.push(validated);
Expand Down Expand Up @@ -148,7 +148,7 @@ const transformBatch = task(
// Subtask: compute statistics from validated records
const computeStatistics = task(
{ name: "computeStatistics", retry },
function computeStatistics(validRecords: ValidatedRecord[]) {
function computeStatistics(_ctx: TaskContext, validRecords: ValidatedRecord[]) {
console.log(`[LOAD] Computing statistics for ${validRecords.length} records`);

if (validRecords.length === 0) {
Expand Down Expand Up @@ -191,24 +191,24 @@ const computeStatistics = task(
// Root task: orchestrates the full ETL pipeline
task(
{ name: "runEtlPipeline", retry, timeoutSeconds: 300 },
async function runEtlPipeline(sourceFile: string) {
async function runEtlPipeline(ctx: TaskContext, sourceFile: string) {
console.log("=".repeat(80));
console.log("[PIPELINE] Starting ETL Pipeline");
console.log(`[PIPELINE] Source: ${sourceFile}`);
console.log("=".repeat(80));

console.log("[PIPELINE] Stage 1/3: EXTRACT");
const rawRecords = await extractCsvData(sourceFile);
const rawRecords = await ctx.step(extractCsvData, sourceFile);
console.log(`[PIPELINE] Extracted ${rawRecords.length} records`);

console.log("[PIPELINE] Stage 2/3: TRANSFORM");
const transformResult = await transformBatch(rawRecords);
const transformResult = await ctx.step(transformBatch, rawRecords);
console.log(
`[PIPELINE] Transformation complete: ${(transformResult.success_rate * 100).toFixed(1)}% success rate`,
);

console.log("[PIPELINE] Stage 3/3: LOAD");
const statistics = await computeStatistics(transformResult.valid_records);
const statistics = await ctx.step(computeStatistics, transformResult.valid_records);
console.log("[PIPELINE] Statistics computed");

const pipelineResult = {
Expand Down
27 changes: 14 additions & 13 deletions file-analyzer/workflow-service/src/main.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import "dotenv/config";
import { task } from "@renderinc/sdk/workflows";
import { task, type TaskContext } from "@renderinc/sdk/workflows";

interface ParsedData {
success: boolean;
Expand All @@ -19,7 +19,7 @@ const retry = {
// Subtask: parse CSV content into structured data
const parseCsvData = task(
{ name: "parseCsvData", retry },
function parseCsvData(fileContent: string): ParsedData {
function parseCsvData(_ctx: TaskContext, fileContent: string): ParsedData {
console.log("[PARSE] Starting CSV parsing");

try {
Expand Down Expand Up @@ -56,7 +56,7 @@ const parseCsvData = task(
// Subtask: calculate statistics from parsed data
const calculateStatistics = task(
{ name: "calculateStatistics", retry },
function calculateStatistics(data: ParsedData) {
function calculateStatistics(_ctx: TaskContext, data: ParsedData) {
console.log("[STATS] Calculating statistics");

if (!data.success || data.rows.length === 0) {
Expand Down Expand Up @@ -111,7 +111,7 @@ const calculateStatistics = task(
// Subtask: identify trends and patterns
const identifyTrends = task(
{ name: "identifyTrends", retry },
function identifyTrends(data: ParsedData) {
function identifyTrends(_ctx: TaskContext, data: ParsedData) {
console.log("[TRENDS] Identifying trends");

if (!data.success || data.rows.length === 0) {
Expand Down Expand Up @@ -164,6 +164,7 @@ const identifyTrends = task(
const generateInsights = task(
{ name: "generateInsights", retry },
async function generateInsights(
_ctx: TaskContext,
stats: { success?: boolean; numeric_columns?: string[]; statistics?: { [col: string]: { avg: number; min: number; max: number; sum: number } } },
trends: { success?: boolean; categorical_columns?: string[]; categorical_analysis?: { [col: string]: { top_5: [string, number][]; distribution: { [key: string]: number } } } },
metadata: ParsedData,
Expand Down Expand Up @@ -214,11 +215,11 @@ const generateInsights = task(
// Root task: orchestrates the full analysis pipeline
task(
{ name: "analyzeFile", retry, timeoutSeconds: 300 },
async function analyzeFile(fileContent: string) {
async function analyzeFile(ctx: TaskContext, fileContent: string) {
console.log("[ANALYZE_FILE] Starting file analysis pipeline");

console.log("[ANALYZE_FILE] Stage 1: Parsing CSV data");
const parsedData = await parseCsvData(fileContent);
const parsedData = await ctx.step(parseCsvData, fileContent);

if (!parsedData.success) {
console.error("[ANALYZE_FILE] Failed to parse CSV data");
Expand All @@ -227,14 +228,14 @@ task(

console.log(`[ANALYZE_FILE] Parsed ${parsedData.row_count} rows`);

console.log("[ANALYZE_FILE] Stage 2: Calculating statistics");
const stats = await calculateStatistics(parsedData);
console.log("[ANALYZE_FILE] Stage 2: Calculating statistics and identifying trends");
const [stats, trends] = await Promise.all([
ctx.step(calculateStatistics, parsedData),
ctx.step(identifyTrends, parsedData),
]);

console.log("[ANALYZE_FILE] Stage 3: Identifying trends");
const trends = await identifyTrends(parsedData);

console.log("[ANALYZE_FILE] Stage 4: Generating insights");
const insights = await generateInsights(stats, trends, parsedData);
console.log("[ANALYZE_FILE] Stage 3: Generating insights");
const insights = await ctx.step(generateInsights, stats, trends, parsedData);

console.log("[ANALYZE_FILE] Analysis pipeline completed successfully");

Expand Down
Loading