Provide a scalable, flexible pipeline to ingest on-chain & off-chain wallet data into ClickHouse in a canonical, queryable form that supports semi-structured/nested payloads. The system supports multiple deployment patterns (one-time, scheduled, event-based) and integrates with blockchain indexing solutions (HyperIndexer, Substreams) for efficient on-chain data extraction.
graph TB
subgraph "Data Sources"
A1[On-Chain Sources]
A2[Off-Chain Sources]
A3[Social Platforms]
A4[DeFi Protocols]
end
subgraph "Indexing Layer"
B1[HyperIndexer]
B2[HyperSync]
B3[HyperRPC]
B4[API Integrations]
end
subgraph "Data Pipe Runtime"
C1[Dagster Orchestrator]
C2[Adapter Registry]
C3[Runtime Engine]
C4[Address Set Manager]
end
subgraph "Streaming Layer"
D1[Kafka/MSK]
D2[Stream Processors]
D3[Transformers]
end
subgraph "Storage Layer"
E1[ClickHouse]
E2[Adapter Registry Table]
E3[User Attribute Tables]
E4[Materialized Views]
end
subgraph "Infrastructure"
F1[AWS CDK]
F2[EKS Cluster]
F3[Monitoring]
F4[Security]
end
A1 --> B1
A1 --> B2
A1 --> B3
A2 --> B4
A3 --> B4
A4 --> B1
B1 --> C1
B2 --> C1
B3 --> C1
B4 --> C1
C1 --> C2
C1 --> C3
C1 --> C4
C3 --> D1
D1 --> D2
D2 --> D3
D3 --> E1
E1 --> E2
E1 --> E3
E1 --> E4
F1 --> F2
F2 --> C1
F2 --> D2
F2 --> E1
sequenceDiagram
participant AS as Address Set Manager
participant AR as Adapter Registry
participant DR as Data Pipe Runtime
participant HI as HyperIndexer
participant K as Kafka
participant SP as Stream Processor
participant CH as ClickHouse
AS->>AR: Register new addresses
AR->>DR: Trigger adapter execution
DR->>HI: Configure indexing for addresses
HI->>HI: Extract blockchain data
HI->>K: Publish events
K->>SP: Stream events
SP->>SP: Transform & normalize
SP->>CH: Insert user attributes
CH->>AR: Update adapter status
graph LR
subgraph "Dagster Assets"
A1[Address Set Asset]
A2[Adapter Config Asset]
A3[Indexing Job Asset]
A4[Data Quality Asset]
end
subgraph "Dagster Jobs"
J1[On-Chain Indexing Job]
J2[Off-Chain Sync Job]
J3[Data Transformation Job]
J4[Quality Check Job]
end
subgraph "Dagster Schedules"
S1[Daily Sync Schedule]
S2[Hourly Index Schedule]
S3[Weekly Backfill Schedule]
end
subgraph "External Systems"
E1[HyperIndexer API]
E2[Kafka Topics]
E3[ClickHouse DB]
E4[Monitoring APIs]
end
A1 --> J1
A2 --> J2
A3 --> J3
A4 --> J4
S1 --> J2
S2 --> J1
S3 --> J3
J1 --> E1
J2 --> E2
J3 --> E3
J4 --> E4
Dagster is ideal for this architecture because it provides:
- Asset-Centric Orchestration: Each adapter, address set, and data transformation becomes a trackable asset
- Dynamic Pipeline Generation: Pipelines can be generated based on adapter configurations
- Data Quality & Lineage: Built-in data quality checks and lineage tracking
- Flexible Scheduling: Supports all deployment types (one-time, scheduled, event-driven)
- Observability: Rich UI for monitoring pipeline execution and data quality
graph TB
subgraph "Dagster Core"
DS[Dagster Server]
DU[Dagster UI]
DW[Dagster Workspace]
end
subgraph "Assets & Jobs"
A1[Address Set Assets]
A2[Adapter Config Assets]
A3[Data Quality Assets]
J1[Indexing Jobs]
J2[Sync Jobs]
J3[Transform Jobs]
end
subgraph "External Integrations"
E1[HyperIndexer API]
E2[Kafka Topics]
E3[ClickHouse DB]
E4[Monitoring Systems]
end
DS --> A1
DS --> A2
DS --> A3
DS --> J1
DS --> J2
DS --> J3
A1 --> E1
A2 --> E2
A3 --> E3
J1 --> E1
J2 --> E2
J3 --> E3
The Data Pipe Runtime is a dynamic execution engine that:
- Adapter Discovery: Discovers and loads adapter implementations
- Dynamic Execution: Executes adapters based on configuration
- Resource Management: Manages compute resources and scaling
- Error Handling: Provides robust error handling and retry mechanisms
graph TB
subgraph "Runtime Engine"
RE[Runtime Engine Core]
AM[Adapter Manager]
RM[Resource Manager]
EM[Execution Manager]
CM[Config Manager]
end
subgraph "Adapter Runtime"
AR1[On-Chain Adapter Runtime]
AR2[Off-Chain Adapter Runtime]
AR3[Hybrid Adapter Runtime]
end
subgraph "Indexing Integration"
II1[HyperIndexer Client]
II2[Substreams Client]
II3[RPC Client]
II4[API Client]
end
subgraph "Execution Context"
EC1[Kubernetes Pods]
EC2[Lambda Functions]
EC3[ECS Tasks]
end
RE --> AM
RE --> RM
RE --> EM
RE --> CM
AM --> AR1
AM --> AR2
AM --> AR3
AR1 --> II1
AR1 --> II2
AR1 --> II3
AR2 --> II4
AR3 --> II1
AR3 --> II4
EM --> EC1
EM --> EC2
EM --> EC3
The system uses a plugin-based architecture for adapter runtime selection:
# Adapter Runtime Interface
class AdapterRuntime(ABC):
@abstractmethod
def execute(self, config: AdapterConfig, context: ExecutionContext) -> ExecutionResult:
pass
@abstractmethod
def validate_config(self, config: AdapterConfig) -> ValidationResult:
pass
@abstractmethod
def get_health_status(self) -> HealthStatus:
pass
# Runtime Selection Logic
class RuntimeSelector:
def select_runtime(self, adapter_id: str, config: AdapterConfig) -> AdapterRuntime:
adapter_info = self.adapter_registry.get_adapter(adapter_id)
if adapter_info.source_type == "onchain":
if adapter_info.indexing_strategy == "hyperindexer":
return HyperIndexerRuntime()
elif adapter_info.indexing_strategy == "substreams":
return SubstreamsRuntime()
else:
return RPCRuntime()
elif adapter_info.source_type == "offchain":
return APIRuntime()
else: # hybrid
return HybridRuntime()Based on the HyperIndexer Python examples, here's how to integrate HyperIndexer for Suite's use cases:
from hypersync import HypersyncClient, Config, Query, FieldSelection, LogField, TransactionField
import asyncio
from typing import List, Dict, Any
class SuiteHyperIndexerClient:
def __init__(self, network: str = "ethereum"):
"""
Initialize HyperIndexer client for Suite's wallet data extraction
Args:
network: Blockchain network (ethereum, polygon, arbitrum, etc.)
"""
self.network = network
self.client = HypersyncClient(network)
async def extract_wallet_transactions(
self,
addresses: List[str],
start_block: int,
end_block: int,
block_range_days: int = 30
) -> List[Dict[str, Any]]:
"""
Extract transactions for wallet addresses in custom block range
Args:
addresses: List of wallet addresses to analyze
start_block: Starting block number
end_block: Ending block number
block_range_days: Number of days for block range (MVP: 30 days)
Returns:
List of transaction data for analytics
"""
# Configure query for transaction extraction
config = Config(
url=f"https://{self.network}.hypersync.xyz",
bearer_token="your_token_here"
)
query = Query(
from_block=start_block,
to_block=end_block,
logs=[
{
"address": addresses, # Filter by wallet addresses
"topics": [
"0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef" # Transfer event
]
}
],
transactions=[
{
"from": addresses, # Transactions from these addresses
"to": addresses # Transactions to these addresses
}
],
field_selection=FieldSelection(
log=[
LogField.ADDRESS,
LogField.TOPICS,
LogField.DATA,
LogField.BLOCK_NUMBER,
LogField.TRANSACTION_HASH,
LogField.LOG_INDEX
],
transaction=[
TransactionField.HASH,
TransactionField.FROM,
TransactionField.TO,
TransactionField.VALUE,
TransactionField.GAS_USED,
TransactionField.GAS_PRICE,
TransactionField.BLOCK_NUMBER,
TransactionField.BLOCK_HASH
]
)
)
# Execute query and collect results
transactions = []
async for data in self.client.stream(query, config):
for tx in data.transactions:
transactions.append({
"wallet_address": tx.from_address if tx.from_address in addresses else tx.to_address,
"transaction_hash": tx.hash,
"from": tx.from_address,
"to": tx.to_address,
"value": tx.value,
"gas_used": tx.gas_used,
"gas_price": tx.gas_price,
"block_number": tx.block_number,
"block_hash": tx.block_hash,
"timestamp": data.block.timestamp
})
return transactions
async def extract_contract_interactions(
self,
addresses: List[str],
contracts: List[str],
start_block: int,
end_block: int
) -> List[Dict[str, Any]]:
"""
Extract contract interactions for Suite's contract-based address import
Args:
addresses: Wallet addresses to analyze
contracts: Contract addresses to monitor
start_block: Starting block number
end_block: Ending block number
Returns:
List of contract interaction data
"""
query = Query(
from_block=start_block,
to_block=end_block,
logs=[
{
"address": contracts, # Monitor specific contracts
"topics": [
"0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef", # Transfer
"0x8c5be1e5ebec7d5bd14f71427d1e84f3dd0314c0f7b2291e5b200ac8c7c3b925" # Approval
]
}
],
transactions=[
{
"from": addresses,
"to": contracts
}
],
field_selection=FieldSelection(
log=[
LogField.ADDRESS,
LogField.TOPICS,
LogField.DATA,
LogField.BLOCK_NUMBER,
LogField.TRANSACTION_HASH
],
transaction=[
TransactionField.HASH,
TransactionField.FROM,
TransactionField.TO,
TransactionField.VALUE,
TransactionField.GAS_USED
]
)
)
interactions = []
async for data in self.client.stream(query):
for log in data.logs:
# Extract wallet addresses from contract interactions
if log.topics and len(log.topics) > 1:
from_address = "0x" + log.topics[1].hex()[-40:]
to_address = "0x" + log.topics[2].hex()[-40:] if len(log.topics) > 2 else None
interactions.append({
"contract_address": log.address,
"wallet_address": from_address,
"interaction_type": "transfer",
"block_number": log.block_number,
"transaction_hash": log.transaction_hash,
"log_data": log.data
})
return interactions
async def extract_nft_activity(
self,
addresses: List[str],
nft_contracts: List[str],
start_block: int,
end_block: int
) -> List[Dict[str, Any]]:
"""
Extract NFT activity for Suite's NFT collection analysis
Args:
addresses: Wallet addresses to analyze
nft_contracts: NFT contract addresses
start_block: Starting block number
end_block: Ending block number
Returns:
List of NFT activity data
"""
query = Query(
from_block=start_block,
to_block=end_block,
logs=[
{
"address": nft_contracts,
"topics": [
"0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef" # Transfer
]
}
],
field_selection=FieldSelection(
log=[
LogField.ADDRESS,
LogField.TOPICS,
LogField.DATA,
LogField.BLOCK_NUMBER,
LogField.TRANSACTION_HASH
]
)
)
nft_activity = []
async for data in self.client.stream(query):
for log in data.logs:
if log.topics and len(log.topics) >= 3:
from_address = "0x" + log.topics[1].hex()[-40:]
to_address = "0x" + log.topics[2].hex()[-40:]
token_id = int(log.topics[3].hex(), 16)
# Check if any of our target addresses are involved
if from_address in addresses or to_address in addresses:
nft_activity.append({
"nft_contract": log.address,
"wallet_address": from_address if from_address in addresses else to_address,
"token_id": token_id,
"action": "transfer_out" if from_address in addresses else "transfer_in",
"block_number": log.block_number,
"transaction_hash": log.transaction_hash
})
return nft_activity
# Suite Adapter Runtime Implementation
class SuiteHyperIndexerRuntime(BaseAdapterRuntime):
def __init__(self, adapter_id: str, config: AdapterConfig):
super().__init__(adapter_id, config)
self.hyperindexer_client = SuiteHyperIndexerClient(
network=config.blockchain_network
)
async def extract_data(self, context: ExecutionContext) -> List[RawData]:
"""Extract data using HyperIndexer based on Suite's use cases"""
# Get addresses from address set
addresses = self.get_addresses_from_set(context.address_set_id)
# Calculate block range (MVP: 30 days)
end_block = await self.hyperindexer_client.client.get_latest_block_number()
start_block = end_block - (30 * 24 * 60 * 60 // 12) # Approximate blocks for 30 days
raw_data = []
if self.config.data_source == "contract_events":
# Extract contract interactions
interactions = await self.hyperindexer_client.extract_contract_interactions(
addresses=addresses,
contracts=self.config.contracts,
start_block=start_block,
end_block=end_block
)
for interaction in interactions:
raw_data.append(RawData(
wallet_address=interaction["wallet_address"],
data=interaction,
timestamp=datetime.utcnow(),
source="hyperindexer_contract"
))
elif self.config.data_source == "nft_marketplace":
# Extract NFT activity
nft_activity = await self.hyperindexer_client.extract_nft_activity(
addresses=addresses,
nft_contracts=self.config.collections,
start_block=start_block,
end_block=end_block
)
for activity in nft_activity:
raw_data.append(RawData(
wallet_address=activity["wallet_address"],
data=activity,
timestamp=datetime.utcnow(),
source="hyperindexer_nft"
))
elif self.config.data_source == "comprehensive_analysis":
# Extract comprehensive transaction data
transactions = await self.hyperindexer_client.extract_wallet_transactions(
addresses=addresses,
start_block=start_block,
end_block=end_block,
block_range_days=30
)
for tx in transactions:
raw_data.append(RawData(
wallet_address=tx["wallet_address"],
data=tx,
timestamp=datetime.utcnow(),
source="hyperindexer_analytics"
))
return raw_datasequenceDiagram
participant AR as Adapter Registry
participant DR as Data Pipe Runtime
participant HI as HyperIndexer
participant K as Kafka
participant SP as Stream Processor
participant CH as ClickHouse
AR->>DR: Trigger adapter execution
DR->>DR: Load adapter configuration
DR->>DR: Select appropriate runtime
DR->>HI: Configure indexing job
HI->>HI: Extract blockchain data
HI->>K: Publish events to topic
K->>SP: Stream events
SP->>SP: Transform & normalize data
SP->>CH: Insert user attributes
CH->>AR: Update adapter status
AR->>DR: Schedule next execution
# Base Adapter Runtime
class BaseAdapterRuntime(AdapterRuntime):
def __init__(self, adapter_id: str, config: AdapterConfig):
self.adapter_id = adapter_id
self.config = config
self.kafka_producer = KafkaProducer()
self.clickhouse_client = ClickHouseClient()
def execute(self, context: ExecutionContext) -> ExecutionResult:
try:
# 1. Validate configuration
validation_result = self.validate_config(self.config)
if not validation_result.is_valid:
return ExecutionResult.error(validation_result.errors)
# 2. Extract data from source
raw_data = self.extract_data(context)
# 3. Transform data
transformed_data = self.transform_data(raw_data)
# 4. Publish to Kafka
self.publish_to_kafka(transformed_data)
# 5. Update status
self.update_adapter_status("completed")
return ExecutionResult.success(len(transformed_data))
except Exception as e:
self.update_adapter_status("error", str(e))
return ExecutionResult.error([str(e)])
def extract_data(self, context: ExecutionContext) -> List[RawData]:
raise NotImplementedError
def transform_data(self, raw_data: List[RawData]) -> List[TransformedData]:
raise NotImplementedError
def publish_to_kafka(self, data: List[TransformedData]):
for item in data:
message = {
"adapter_id": self.adapter_id,
"wallet_address": item.wallet_address,
"payload": item.data,
"occurred_at": item.timestamp,
"event_id": str(uuid.uuid4())
}
self.kafka_producer.send(
topic=f"adapter.{self.adapter_id}.events",
key=item.wallet_address,
value=json.dumps(message)
)
# HyperIndexer Runtime Implementation
class HyperIndexerRuntime(BaseAdapterRuntime):
def __init__(self, adapter_id: str, config: AdapterConfig):
super().__init__(adapter_id, config)
self.hyperindexer_client = HyperIndexerClient(
api_key=config.hyperindexer_api_key,
base_url=config.hyperindexer_base_url
)
def extract_data(self, context: ExecutionContext) -> List[RawData]:
# Get addresses from address set
addresses = self.get_addresses_from_set(context.address_set_id)
# Create indexing job
indexing_config = HyperIndexerConfig(
addresses=addresses,
contracts=self.config.contracts,
events=self.config.events,
start_block=self.config.start_block,
end_block=self.config.end_block
)
job = self.hyperindexer_client.create_indexing_job(indexing_config)
# Stream events from the job
raw_data = []
for event in self.hyperindexer_client.stream_events(job.id):
raw_data.append(RawData(
wallet_address=event.address,
data=event.data,
timestamp=event.timestamp,
source="hyperindexer"
))
return raw_data
def transform_data(self, raw_data: List[RawData]) -> List[TransformedData]:
transformed_data = []
for item in raw_data:
# Flatten nested JSON into key-value pairs
flattened = self.flatten_json(item.data)
for key_path, value in flattened.items():
transformed_data.append(TransformedData(
wallet_address=item.wallet_address,
key_path=key_path,
field_name=key_path.split('.')[-1],
field_type=self.detect_field_type(value),
field_value=str(value),
raw_json=json.dumps(item.data),
occurred_at=item.timestamp
))
return transformed_data
def flatten_json(self, data: dict, prefix: str = "") -> dict:
"""Flatten nested JSON into key-value pairs"""
flattened = {}
for key, value in data.items():
new_key = f"{prefix}.{key}" if prefix else key
if isinstance(value, dict):
flattened.update(self.flatten_json(value, new_key))
elif isinstance(value, list):
for i, item in enumerate(value):
if isinstance(item, dict):
flattened.update(self.flatten_json(item, f"{new_key}.{i}"))
else:
flattened[f"{new_key}.{i}"] = item
else:
flattened[new_key] = value
return flattened
# Off-Chain API Runtime Implementation
class APIRuntime(BaseAdapterRuntime):
def __init__(self, adapter_id: str, config: AdapterConfig):
super().__init__(adapter_id, config)
self.api_client = APIClient(
base_url=config.api_base_url,
api_key=config.api_key,
rate_limit=config.rate_limit
)
def extract_data(self, context: ExecutionContext) -> List[RawData]:
raw_data = []
# Get addresses from address set
addresses = self.get_addresses_from_set(context.address_set_id)
for address in addresses:
try:
# Make API call with rate limiting
response = self.api_client.get_user_data(address)
raw_data.append(RawData(
wallet_address=address,
data=response.json(),
timestamp=datetime.utcnow(),
source="api"
))
# Respect rate limits
time.sleep(self.config.rate_limit_delay)
except Exception as e:
logger.error(f"Failed to fetch data for {address}: {e}")
continue
return raw_datafrom dagster import asset, job, schedule, sensor, AssetExecutionContext
from dagster_aws import S3Resource
from dagster_k8s import K8sJobOp
# Address Set Asset
@asset(
description="Manages the set of wallet addresses to be indexed",
group_name="address_management"
)
def address_set_asset(context: AssetExecutionContext) -> dict:
"""Asset representing the current set of addresses to be indexed"""
address_set_manager = AddressSetManager()
# Get addresses from various sources
addresses = address_set_manager.get_active_addresses()
# Update address set registry
address_set_manager.update_address_set("active_addresses", addresses)
return {
"addresses": addresses,
"count": len(addresses),
"last_updated": datetime.utcnow().isoformat()
}
# Adapter Configuration Asset
@asset(
description="Adapter configuration and metadata",
group_name="adapter_management"
)
def adapter_config_asset(context: AssetExecutionContext) -> dict:
"""Asset representing adapter configurations"""
adapter_registry = AdapterRegistry()
# Get all active adapters
adapters = adapter_registry.get_active_adapters()
return {
"adapters": adapters,
"count": len(adapters),
"last_updated": datetime.utcnow().isoformat()
}
# Data Quality Asset
@asset(
description="Data quality metrics and validation results",
group_name="data_quality"
)
def data_quality_asset(
context: AssetExecutionContext,
address_set: dict,
adapter_config: dict
) -> dict:
"""Asset representing data quality metrics"""
quality_checker = DataQualityChecker()
# Check data quality for each adapter
quality_metrics = {}
for adapter in adapter_config["adapters"]:
metrics = quality_checker.check_adapter_quality(adapter["adapter_id"])
quality_metrics[adapter["adapter_id"]] = metrics
return {
"quality_metrics": quality_metrics,
"overall_score": quality_checker.calculate_overall_score(quality_metrics),
"last_checked": datetime.utcnow().isoformat()
}# Suite Contract-Based Address Import Job
@job(
description="Job for extracting wallet addresses from contract interactions",
tags={"type": "onchain", "source": "contract_events", "suite": "address_import"}
)
def contract_address_import_job():
"""Job that extracts wallet addresses from contract events using HyperIndexer"""
# Get address set
address_set = address_set_asset()
# Get adapter config
adapter_config = adapter_config_asset()
# Run indexing for contract-based adapters
for adapter in adapter_config["adapters"]:
if adapter["data_source"] == "contract_events":
K8sJobOp(
name=f"contract-import-{adapter['adapter_id']}",
image=f"wallet-pipes/{adapter['adapter_id']}:{adapter['adapter_version']}",
env_vars={
"ADAPTER_ID": adapter["adapter_id"],
"ADDRESS_SET_ID": "active_addresses",
"KAFKA_BROKERS": "wallet-data-pipeline:9092",
"BLOCKCHAIN_NETWORK": adapter["blockchain_network"],
"BLOCK_RANGE_DAYS": "30"
}
)()
# Suite NFT Collection Analysis Job
@job(
description="Job for analyzing NFT collections and extracting wallet addresses",
tags={"type": "hybrid", "source": "nft_marketplace", "suite": "nft_analysis"}
)
def nft_collection_analysis_job():
"""Job that analyzes NFT collections using HyperIndexer"""
# Get address set
address_set = address_set_asset()
# Get adapter config
adapter_config = adapter_config_asset()
# Run analysis for NFT collection adapters
for adapter in adapter_config["adapters"]:
if adapter["data_source"] == "nft_marketplace":
K8sJobOp(
name=f"nft-analysis-{adapter['adapter_id']}",
image=f"wallet-pipes/{adapter['adapter_id']}:{adapter['adapter_version']}",
env_vars={
"ADAPTER_ID": adapter["adapter_id"],
"ADDRESS_SET_ID": "active_addresses",
"KAFKA_BROKERS": "wallet-data-pipeline:9092",
"BLOCKCHAIN_NETWORK": adapter["blockchain_network"],
"BLOCK_RANGE_DAYS": "30"
}
)()
# Suite Off-Chain Data Sync Job
@job(
description="Job for syncing off-chain data from Privy, Guild.xyz, etc.",
tags={"type": "offchain", "source": "api", "suite": "offchain_sync"}
)
def suite_offchain_sync_job():
"""Job that syncs off-chain data from Suite's data sources"""
# Get address set
address_set = address_set_asset()
# Get adapter config
adapter_config = adapter_config_asset()
# Run sync for off-chain adapters
for adapter in adapter_config["adapters"]:
if adapter["data_source"] in ["privy_api", "guild_api"]:
K8sJobOp(
name=f"offchain-sync-{adapter['adapter_id']}",
image=f"wallet-pipes/{adapter['adapter_id']}:{adapter['adapter_version']}",
env_vars={
"ADAPTER_ID": adapter["adapter_id"],
"ADDRESS_SET_ID": "active_addresses",
"KAFKA_BROKERS": "wallet-data-pipeline:9092",
"API_BASE_URL": adapter["config"]["api_base_url"],
"API_KEY": adapter["config"]["api_key"]
}
)()
# Suite Generic On-Chain Analytics Job
@job(
description="Job for comprehensive on-chain analytics across all addresses",
tags={"type": "onchain", "source": "comprehensive_analysis", "suite": "analytics"}
)
def suite_generic_analytics_job():
"""Job that performs comprehensive on-chain analytics using HyperIndexer"""
# Get address set
address_set = address_set_asset()
# Get adapter config
adapter_config = adapter_config_asset()
# Run analytics for comprehensive analysis adapters
for adapter in adapter_config["adapters"]:
if adapter["data_source"] == "comprehensive_analysis":
K8sJobOp(
name=f"analytics-{adapter['adapter_id']}",
image=f"wallet-pipes/{adapter['adapter_id']}:{adapter['adapter_version']}",
env_vars={
"ADAPTER_ID": adapter["adapter_id"],
"ADDRESS_SET_ID": "all_active_addresses",
"KAFKA_BROKERS": "wallet-data-pipeline:9092",
"BLOCKCHAIN_NETWORK": adapter["blockchain_network"],
"BLOCK_RANGE_DAYS": "30",
"ANALYSIS_TYPES": ",".join(adapter["config"]["analysis_types"])
}
)()
# Data Transformation Job
@job(
description="Job for transforming and normalizing data",
tags={"type": "transformation", "processing": "stream"}
)
def data_transformation_job():
"""Job that orchestrates data transformation"""
# Get adapter config
adapter_config = adapter_config_asset()
# Run transformation for each adapter
for adapter in adapter_config["adapters"]:
K8sJobOp(
name=f"transform-{adapter['adapter_id']}",
image="wallet-pipes/stream-processor:latest",
env_vars={
"ADAPTER_ID": adapter["adapter_id"],
"KAFKA_BROKERS": "wallet-data-pipeline:9092",
"CLICKHOUSE_URL": "clickhouse://wallet-data-ch:9000"
}
)()# Suite Contract Address Import Schedule (Kafka-triggered)
@schedule(
job=contract_address_import_job,
cron_schedule="0 */6 * * *", # Every 6 hours
description="Regular contract address import using HyperIndexer"
)
def contract_import_schedule(context):
"""Schedule for contract-based address import"""
return {}
# Suite NFT Collection Analysis Schedule (Kafka-triggered)
@schedule(
job=nft_collection_analysis_job,
cron_schedule="0 */4 * * *", # Every 4 hours
description="Regular NFT collection analysis using HyperIndexer"
)
def nft_analysis_schedule(context):
"""Schedule for NFT collection analysis"""
return {}
# Suite Off-Chain Data Sync Schedule (Kafka-triggered)
@schedule(
job=suite_offchain_sync_job,
cron_schedule="0 */2 * * *", # Every 2 hours
description="Regular off-chain data sync from Privy, Guild.xyz, etc."
)
def offchain_sync_schedule(context):
"""Schedule for off-chain data synchronization"""
return {}
# Suite Generic Analytics Schedule (Batch Scheduled)
@schedule(
job=suite_generic_analytics_job,
cron_schedule="0 1 * * *", # Daily at 1 AM
description="Daily comprehensive on-chain analytics using HyperIndexer"
)
def generic_analytics_schedule(context):
"""Schedule for comprehensive on-chain analytics"""
return {}
# Suite Data Transformation Schedule
@schedule(
job=data_transformation_job,
cron_schedule="0 */1 * * *", # Every hour
description="Hourly data transformation and normalization"
)
def data_transformation_schedule(context):
"""Schedule for data transformation"""
return {}-
Pipe Registry & Management
adapter_registrytable stores adapter metadata, deployment configurations, and runtime status- Pipe deployment types: one-time jobs, scheduled tasks, event-driven (Kafka-triggered)
- Dynamic address set management for on-chain indexing
-
Data Sources & Adapters
On-Chain Sources:
- HyperIndexer integration for ultra-fast blockchain data extraction
- HyperSync for high-performance data streaming (1000x faster than RPC)
- HyperRPC for read-only RPC operations
- Custom block range extraction for analytics (MVP: 1-month range)
Off-Chain Sources:
- Social platforms (Twitter/X, GitHub, LinkedIn)
- DeFi protocols (Uniswap, Aave, Compound)
- NFT marketplaces (OpenSea, Blur)
- Traditional finance APIs
- Webhook integrations
-
Streaming Infrastructure: Kafka (AWS MSK or MSK Serverless)
- Durability, backpressure, replay, consumer groups
- Topic partitioning by
wallet_addresshash for locality - Dead letter queues for error handling
-
Stream Processors & Transformers
- Stateless processors for simple transformations
- Stateful processors (Flink/Kafka Streams) for complex aggregations
- Real-time normalization and flattening of nested JSON
- Address set synchronization between on-chain and off-chain sources
-
ClickHouse Data Layer
- Canonical
user_attributetable for normalized data adapter_registrytable for adapter metadata- Materialized views for latest state aggregation
- Auto-table creation and schema evolution
- Canonical
-
Infrastructure & Operations
- AWS CDK for infrastructure provisioning
- Kubernetes (EKS) for container orchestration
- GitOps for adapter deployment and configuration
- Comprehensive monitoring and alerting
The adapter_registry table stores adapter metadata separately from user data, enabling better governance and management.
CREATE TABLE IF NOT EXISTS adapter_registry (
adapter_id String, -- unique identifier
adapter_name String, -- human-readable name
adapter_version String, -- semantic version
schema_version String, -- schema version
table_name String, -- target ClickHouse table
deployment_type Enum8( -- deployment pattern
'one_time' = 1,
'scheduled' = 2,
'event_driven' = 3,
'continuous' = 4
),
source_type Enum8( -- data source type
'onchain' = 1,
'offchain' = 2,
'hybrid' = 3
),
blockchain_network String, -- e.g., "ethereum", "polygon", "arbitrum"
indexing_strategy String, -- e.g., "hyperindexer", "substreams", "rpc"
address_set_strategy String, -- how addresses are managed
config JSON, -- adapter-specific configuration
status Enum8( -- runtime status
'active' = 1,
'paused' = 2,
'error' = 3,
'deprecated' = 4
),
last_run_at DateTime64(3),
next_run_at DateTime64(3),
created_at DateTime64(3),
updated_at DateTime64(3)
) ENGINE = ReplacingMergeTree(updated_at)
ORDER BY (adapter_id, adapter_version);Simplified user data table without adapter metadata duplication.
CREATE TABLE IF NOT EXISTS {table_name} (
wallet_address String, -- canonicalized address
adapter_id String, -- reference to adapter_registry
key_path Array(String), -- e.g., ['profile','social','twitter']
key_path_str String, -- compact: 'profile.social.twitter'
field_name String, -- final field name
field_type String, -- e.g., "string", "number", "boolean", "json"
field_value String, -- stringified value
raw_json String, -- full original JSON (optional)
occurred_at DateTime64(3), -- event timestamp
ingested_at DateTime64(3) -- ingestion timestamp
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(occurred_at)
ORDER BY (wallet_address, adapter_id, key_path_str, field_name, occurred_at);- Use case: Historical data backfill, one-off data migrations
- Trigger: Manual or API call
- Lifecycle: Run once and terminate
- Example: Backfill all Uniswap transactions for a specific address set
- Use case: Periodic data updates, batch processing
- Trigger: Cron schedule or interval-based
- Lifecycle: Recurring execution
- Example: Daily GitHub activity sync, weekly portfolio rebalancing
- Use case: Real-time processing, reactive updates
- Trigger: Kafka message consumption
- Lifecycle: Long-running consumer
- Example: Process new transactions as they appear on-chain
- Use case: Real-time monitoring, live data feeds
- Trigger: Always-on stream processing
- Lifecycle: Persistent service
- Example: Live DeFi position monitoring, social media sentiment tracking
Address Set Management:
{
"address_set_id": "defi_users_v1",
"addresses": ["0xabc...", "0xdef..."],
"update_strategy": "incremental",
"sync_frequency": "5m",
"indexing_config": {
"contracts": [
{
"address": "0xA0b86a33E6441b8C4C8C0C8C0C8C0C8C0C8C0C8C",
"events": ["Transfer", "Swap"],
"methods": ["balanceOf", "totalSupply"]
}
],
"blocks": {
"start_block": 18000000,
"end_block": "latest"
}
}
}HyperIndexer Adapter Pattern:
- Subscribe to address set changes via webhook
- Configure indexing parameters per contract/event
- Handle incremental updates efficiently
- Support multiple blockchain networks
Real-time Event Processing:
# substreams.yaml
modules:
- name: wallet_events
kind: map
inputs:
- source: sf.ethereum.type.v2.Block
output:
type: proto:wallet.WalletEventSubstreams Adapter Features:
- Real-time blockchain event streaming
- Custom event filtering and transformation
- Efficient data extraction for specific addresses
- Support for complex event relationships
Fallback for Custom Queries:
- Direct blockchain RPC calls
- Custom smart contract method calls
- Historical data queries
- Gas optimization strategies
Twitter/X Integration:
{
"adapter_id": "twitter_profile_adapter",
"source_type": "offchain",
"deployment_type": "scheduled",
"config": {
"api_endpoints": ["/2/users/by/username", "/2/users/by"],
"rate_limits": {
"requests_per_15min": 300,
"requests_per_day": 10000
},
"fields": ["username", "name", "description", "public_metrics"],
"update_frequency": "1h"
}
}GitHub Integration:
{
"adapter_id": "github_activity_adapter",
"source_type": "offchain",
"deployment_type": "event_driven",
"config": {
"webhook_events": ["push", "pull_request", "issues"],
"api_endpoints": ["/user", "/user/repos", "/user/events"],
"fields": ["login", "name", "bio", "public_repos", "followers"],
"repository_filter": "wallet-connected"
}
}Uniswap Integration:
{
"adapter_id": "uniswap_positions_adapter",
"source_type": "hybrid",
"deployment_type": "continuous",
"config": {
"contracts": {
"v3_factory": "0x1F98431c8aD98523631AE4a59f267346ea31F984",
"v3_positions": "0xC36442b4a4522E871399CD717aBDD847Ab11FE88"
},
"events": ["IncreaseLiquidity", "DecreaseLiquidity", "Collect"],
"methods": ["positions", "tokenOfOwnerByIndex"],
"update_frequency": "30s"
}
}OpenSea Integration:
{
"adapter_id": "opensea_collections_adapter",
"source_type": "offchain",
"deployment_type": "scheduled",
"config": {
"api_version": "v1",
"endpoints": ["/assets", "/collections", "/events"],
"fields": ["name", "description", "image_url", "floor_price"],
"update_frequency": "6h"
}
}Address Set Registry:
CREATE TABLE IF NOT EXISTS address_sets (
set_id String,
set_name String,
addresses Array(String),
source_type Enum8('manual' = 1, 'discovered' = 2, 'derived' = 3),
discovery_rules JSON,
last_updated DateTime64(3),
created_at DateTime64(3)
) ENGINE = ReplacingMergeTree(last_updated)
ORDER BY set_id;Address Discovery Strategies:
- Manual: Curated address lists
- Discovered: Addresses found through transaction analysis
- Derived: Addresses computed from other data sources
Selective Indexing:
- Index only addresses in active sets
- Implement address importance scoring
- Use tiered indexing (high-value addresses get more frequent updates)
- Archive inactive addresses to cold storage
Storage Optimization:
-- Archive old data for inactive addresses
CREATE TABLE address_archive AS address_sets
ENGINE = MergeTree()
ORDER BY set_id
TTL last_updated + INTERVAL 90 DAY;{
"adapter_id": "uniswap_v3_positions",
"adapter_name": "Uniswap V3 Position Tracker",
"adapter_version": "2.1.0",
"schema_version": "2025-01-15:v2",
"table_name": "defi_positions",
"deployment_type": "continuous",
"source_type": "hybrid",
"blockchain_network": "ethereum",
"indexing_strategy": "hyperindexer",
"address_set_strategy": "dynamic_discovery",
"config": {
"contracts": {
"positions_manager": "0xC36442b4a4522E871399CD717aBDD847Ab11FE88",
"factory": "0x1F98431c8aD98523631AE4a59f267346ea31F984"
},
"events": ["IncreaseLiquidity", "DecreaseLiquidity", "Collect"],
"update_frequency": "30s",
"batch_size": 1000,
"retry_policy": {
"max_retries": 3,
"backoff_multiplier": 2
}
},
"monitoring": {
"health_check_endpoint": "/health",
"metrics_endpoint": "/metrics",
"alert_thresholds": {
"error_rate": 0.05,
"latency_p99": 5000
}
}
}Address Set Registry → HyperIndexer → Kafka Topic → Stream Processor → ClickHouse
- Address Discovery: Monitor transactions, social connections, DeFi interactions
- Indexing: HyperIndexer extract relevant blockchain data
- Streaming: Data flows through Kafka with wallet-based partitioning
- Processing: Stream processors normalize and flatten data
- Storage: Data stored in ClickHouse with proper indexing
External APIs → Adapter Services → Kafka Topic → Stream Processor → ClickHouse
- API Integration: Adapters fetch data from external services
- Rate Limiting: Implement proper rate limiting and backoff strategies
- Streaming: Normalized data flows through Kafka
- Processing: Stream processors handle data transformation
- Storage: Data stored in ClickHouse with deduplication
// Core infrastructure components
const mskCluster = new kafka.CfnCluster(this, "WalletDataMSK", {
clusterName: "wallet-data-pipeline",
kafkaVersion: "2.8.1",
numberOfBrokerNodes: 3,
brokerNodeGroupInfo: {
instanceType: "kafka.m5.large",
storageInfo: {
ebsStorageInfo: {
volumeSize: 100,
},
},
},
});
const eksCluster = new eks.Cluster(this, "WalletDataEKS", {
version: eks.KubernetesVersion.V1_28,
defaultCapacity: 2,
defaultCapacityInstance: ec2.InstanceType.of(
ec2.InstanceClass.M5,
ec2.InstanceSize.LARGE
),
});
const clickHouseCluster = new clickhouse.CfnCluster(
this,
"WalletDataClickHouse",
{
clusterName: "wallet-data-ch",
nodeType: "clickhouse.x1.medium",
numberOfNodes: 3,
}
);Adapter Deployment:
apiVersion: apps/v1
kind: Deployment
metadata:
name: uniswap-adapter
spec:
replicas: 2
selector:
matchLabels:
app: uniswap-adapter
template:
metadata:
labels:
app: uniswap-adapter
spec:
containers:
- name: adapter
image: wallet-pipes/uniswap-adapter:v2.1.0
env:
- name: KAFKA_BROKERS
value: "wallet-data-pipeline:9092"
- name: CLICKHOUSE_URL
valueFrom:
secretKeyRef:
name: clickhouse-secret
key: url
resources:
requests:
memory: "512Mi"
cpu: "250m"
limits:
memory: "1Gi"
cpu: "500m"Adapter Metrics:
- Processing rate (events/second)
- Error rate and types
- Latency percentiles (p50, p95, p99)
- Address set size and update frequency
- Data freshness (time since last update)
Infrastructure Metrics:
- Kafka consumer lag
- ClickHouse query performance
- Resource utilization (CPU, memory, disk)
- Network throughput and latency
# Prometheus alerting rules
groups:
- name: wallet-data-pipeline
rules:
- alert: HighErrorRate
expr: rate(adapter_errors_total[5m]) > 0.1
for: 2m
labels:
severity: warning
annotations:
summary: "High error rate detected"
- alert: KafkaConsumerLag
expr: kafka_consumer_lag_sum > 10000
for: 5m
labels:
severity: critical
annotations:
summary: "Kafka consumer lag is high"PII Handling:
- Field-level encryption for sensitive data
- Data masking for non-production environments
- Audit logging for data access
- GDPR compliance for EU users
Access Control:
- Role-based access control (RBAC)
- API key management for external services
- Network segmentation and firewall rules
- Secrets management with AWS Secrets Manager
Address Validation:
- Checksum validation for Ethereum addresses
- Multi-signature wallet support
- Smart contract address verification
- Transaction signature validation
Data Lifecycle Management:
-- Hot data (last 30 days)
CREATE TABLE user_attribute_hot AS user_attribute
ENGINE = MergeTree()
ORDER BY (wallet_address, adapter_id, key_path_str, field_name, occurred_at)
TTL occurred_at + INTERVAL 30 DAY;
-- Warm data (30-90 days, compressed)
CREATE TABLE user_attribute_warm AS user_attribute
ENGINE = MergeTree()
ORDER BY (wallet_address, adapter_id, key_path_str, field_name, occurred_at)
TTL occurred_at + INTERVAL 90 DAY
SETTINGS storage_policy = 'warm_storage';
-- Cold data (90+ days, archived to S3)
CREATE TABLE user_attribute_cold AS user_attribute
ENGINE = S3('s3://wallet-data-archive/user_attribute/', 'Parquet')
ORDER BY (wallet_address, adapter_id, key_path_str, field_name, occurred_at);Resource Scaling:
- Horizontal pod autoscaling based on Kafka lag
- Vertical pod autoscaling for resource optimization
- Spot instances for non-critical workloads
- Reserved instances for predictable workloads
- Set up basic infrastructure (MSK, EKS, ClickHouse)
- Implement adapter registry and user attribute tables
- Create first on-chain adapter (HyperIndexer integration)
- Basic monitoring and alerting
- Implement all deployment types (one-time, scheduled, event-driven)
- Add off-chain adapters (Twitter, GitHub)
- Address set management and discovery
- Enhanced monitoring and cost optimization
- Multi-blockchain support
- Advanced analytics and materialized views
- Comprehensive security and compliance
- Performance optimization and cost reduction
- Multi-region deployment
- Disaster recovery and backup strategies
- Advanced monitoring and alerting
- Documentation and operational runbooks
SELECT
wallet_address,
key_path_str,
field_value,
occurred_at
FROM user_attribute_latest
WHERE wallet_address = '0xabc...'
AND adapter_id IN (
SELECT adapter_id
FROM adapter_registry
WHERE source_type = 'offchain'
)
ORDER BY occurred_at DESC;SELECT
wallet_address,
key_path_str,
field_value,
occurred_at
FROM defi_positions
WHERE wallet_address IN (
SELECT addresses
FROM address_sets
WHERE set_name = 'defi_users_v1'
)
AND key_path_str LIKE 'position.%'
AND occurred_at >= now() - INTERVAL 7 DAY
ORDER BY wallet_address, occurred_at DESC;SELECT
wallet_address,
blockchain_network,
count(*) as activity_count
FROM user_attribute ua
JOIN adapter_registry ar ON ua.adapter_id = ar.adapter_id
WHERE ua.occurred_at >= now() - INTERVAL 24 HOUR
AND ar.source_type = 'onchain'
GROUP BY wallet_address, blockchain_network
ORDER BY activity_count DESC;The system supports real-time/streaming address processing for immediate persona updates. Use these commands to test the streaming pipeline:
# Check Kafka connectivity and streaming setup
just stream-check
# Enqueue a single test address
just stream-test
# Enqueue a specific address
just stream-test 0x1234567890abcdef1234567890abcdef12345678
# Enqueue multiple test addresses
just stream-batch 10
# Monitor Kafka topic in real-time
just stream-monitor-
Start Infrastructure
just up
-
Enable Streaming Sensors (in Dagster UI)
- Go to Automation → Sensors
- Start
new_address_sensor - Start
organization_import_sensor
-
Enqueue Test Address
just stream-test 0xYourTestAddress
-
Monitor Processing
- Watch Dagster UI for sensor ticks and adapter runs
- Check ClickHouse for processed data:
SELECT * FROM user_attribute WHERE wallet_address = '0xYourTestAddress' ORDER BY ingested_at DESC;
- STREAMING_TESTING.md - Complete testing guide
- STREAMING_IMPLEMENTATION_GUIDE.md - Implementation details
- ADDRESS_SET_ARCHITECTURE.md - Architecture overview
- Separate Concerns: Use
adapter_registrytable for metadata,user_attributefor data - Flexible Deployment: Support one-time, scheduled, event-driven, and continuous deployment patterns
- On-Chain Integration: Leverage HyperIndexer and Substreams for efficient blockchain data extraction
- Off-Chain Sources: Comprehensive support for social platforms, DeFi protocols, and NFT marketplaces
- Address Management: Dynamic address set discovery and cost-optimized indexing
- Infrastructure: AWS CDK + EKS + MSK + ClickHouse for scalable, managed infrastructure
- Monitoring: Comprehensive observability with Prometheus, Grafana, and alerting
- Security: Field-level encryption, RBAC, and compliance with data privacy regulations
- Cost Optimization: Tiered storage, selective indexing, and resource scaling strategies