← AlliumCONTENT HISTORYWHAT CHANGED · RULE-BASED ANALYSIS
Update to Allium
Snapshot Sep 30, 2026 · 23:01 UTC · version 1.0.0
Collection source: not recorded for this historical snapshot.
First saved snapshot
No earlier snapshot is available to establish a change.
Compare saved observations
Download comparison JSONFull technical diff · 0 changed fields
Full snapshot data
{
"name": "beam-pipelines",
"description": "**Required for Beam pipeline tools: create_beam_config, deploy_beam_pipeline, etc.**\n**Beam** is Allium's custom real-time data pipeline product, built on top of\n[Allium Datastreams](https://app.allium.so/build/datastreams). It lets you tap into\nany of Allium's real-time blockchain data streams — covering 80+ chains — and apply\nyour own filters and JavaScript transformations before delivering the results to\nKafka, SNS, or other destinations.\n\nCommon use cases: real-time alerts on contract events, custom data feeds filtered\nto specific wallets or protocols, and streaming enriched data to your own pipeline.\n\nBrowse available data streams at [app.allium.so/build/datastreams](https://app.allium.so/build/datastreams).\nTo get access or learn more, reach out to support@allium.so.\n\nRead this skill for detailed configuration guides and examples.\n",
"included_files": [],
"skill_md_contents": "---\nname: beam-pipelines\ndescription: |\n **Required for Beam pipeline tools: create_beam_config, deploy_beam_pipeline, etc.**\n **Beam** is Allium's custom real-time data pipeline product, built on top of\n [Allium Datastreams](https://app.allium.so/build/datastreams). It lets you tap into\n any of Allium's real-time blockchain data streams — covering 80+ chains — and apply\n your own filters and JavaScript transformations before delivering the results to\n Kafka, SNS, or other destinations.\n\n Common use cases: real-time alerts on contract events, custom data feeds filtered\n to specific wallets or protocols, and streaming enriched data to your own pipeline.\n\n Browse available data streams at [app.allium.so/build/datastreams](https://app.allium.so/build/datastreams).\n To get access or learn more, reach out to support@allium.so.\n\n Read this skill for detailed configuration guides and examples.\ngate: app__custom_transforms\n---\n\n# Beam Pipeline Configuration Guide\n\n## What is Beam?\n\nBeam is Allium's custom real-time data pipeline product, built on top of [Allium Datastreams](https://app.allium.so/build/datastreams). It lets you tap into any of Allium's real-time blockchain data streams — covering 80+ chains — and apply your own filters and transformations before delivering the results to your preferred destination.\n\n```text\nSource (any Allium Datastream) → Transforms (filter/transform) → Sinks (Kafka, SNS, and more)\n```\n\nWith Beam, you can:\n\n- **Filter** high-volume blockchain streams down to just the events you care about (specific contracts, addresses, event signatures)\n- **Transform** data in-flight using JavaScript to parse, enrich, or reshape records before delivery\n- **Deliver** processed data to Kafka topics or SNS (with more sink types coming soon) for consumption by your applications\n\nBrowse the full catalog of available data streams at [app.allium.so/build/datastreams](https://app.allium.so/build/datastreams).\n\n**Common use cases:**\n\n- Real-time alerts on contract events (e.g., large token transfers, liquidations)\n- Custom data feeds filtered to specific wallet addresses or protocols\n- Live monitoring and anomaly detection for DeFi activity\n- Streaming enriched trade data to your own analytics pipeline\n\n**Interested in Beam?** Reach out to <support@allium.so> to get access.\n\n## Quick Start Workflow\n\n1. **Create config** → `create_beam_config`\n2. **Deploy pipeline** → `deploy_beam_pipeline`\n3. **Verify deployment** → `get_beam_deployment_stats`\n\nAlways check deployment stats after deploying to confirm workers are healthy.\n\nAfter deploying, you'll receive:\n\n- **Kafka connection credentials** (bootstrap server, username, password)\n- **Consumer code snippets** in Python and TypeScript, ready to copy-paste\n- A **UI link** to manage your pipeline at `app.allium.so/build/beam/{config_id}`\n\n**Updating a pipeline:** Just call `deploy_beam_pipeline` again after updating the config.\nDeployment is idempotent and applies changes in place with zero downtime.\nDo NOT teardown before redeploying—that causes unnecessary downtime.\n\n## Configuration Reference\n\n### Source\n\nBeam sources connect to Allium's Datastreams. You select a chain and entity type to stream.\n\n```json\n{\n \"chain\": \"polygon\",\n \"entity\": \"log\",\n \"is_zerolag\": false\n}\n```\n\n| Field | Required | Description |\n| ------------ | -------- | ---------------------------------------- |\n| `chain` | Yes | Blockchain to source data from |\n| `entity` | Yes | Entity type to stream (see table below) |\n| `is_zerolag` | No | Low-latency source mode (default: false) |\n\n#### Supported Chains & Entities\n\nCheck the [Datastreams catalog](https://app.allium.so/build/datastreams) for the latest availability, or reach out to <support@allium.so> to request a specific chain or entity.\n\n### Transforms\n\nTwo transform types are available. You can chain multiple transforms together — data flows through them in order.\n\n#### redis_set_filter\n\nFilter data by matching field values against a set. Only records whose extracted field value exists in your set will pass through.\n\n```json\n{\n \"type\": \"redis_set_filter\",\n \"set_values\": [\"0x1234...\", \"0x5678...\"],\n \"filter_expr\": \"root = this.address\"\n}\n```\n\n| Field | Description |\n| ------------- | ------------------------------------------------------------ |\n| `set_values` | Array of values to match against |\n| `filter_expr` | Bloblang expression to extract the field value for filtering |\n\n**Important:** When filtering by addresses, labels, or symbols, use **lowercase** values.\nThis is how these values are normalized in our system.\n\n**Bloblang `filter_expr` examples:**\n\n```bloblang\nroot = this.address # Filter by contract address\nroot = this.topic0 # Filter by first topic (event signature)\nroot = this.from_address # Filter by sender address\n```\n\n**Log entity field reference:**\n\n| Field | Description |\n| -------------- | ------------------------------------------ |\n| `address` | Contract address that emitted the log |\n| `topic0` | First topic (usually event signature hash) |\n| `topic1` | Second topic (first indexed parameter) |\n| `topic2` | Third topic (second indexed parameter) |\n| `topic3` | Fourth topic (third indexed parameter) |\n| `from_address` | Transaction sender |\n| `to_address` | Transaction recipient |\n| `data` | Non-indexed event data |\n\n#### v8\n\nTransform data using JavaScript. Your function receives each record and can modify, enrich, or reshape it. Return `null` to drop a record.\n\n```json\n{\n \"type\": \"v8\",\n \"script\": \"function transform(record) { return record; }\"\n}\n```\n\n| Field | Description |\n| -------- | ------------------------------------------ |\n| `script` | JavaScript code that processes each record |\n\n### Sinks\n\nSinks define where your processed data is delivered.\n\n```json\n{\n \"type\": \"kafka\",\n \"name\": \"my-output-topic\"\n}\n```\n\n| Field | Description |\n| ------ | ------------------------------------------------------ |\n| `type` | Output type: `kafka` or `sns` (more sinks coming soon) |\n| `name` | Topic name suffix for the output |\n\n**Kafka sinks:** After deployment, you'll receive connection credentials and ready-to-use consumer code snippets in Python and TypeScript.\n\n**SNS sinks:** Delivers data to an SNS topic. Reach out to <support@allium.so> if you need help setting up SNS delivery or are interested in other sink types.\n\n## Quick Start Examples\n\n### Example 1: Filter Logs by Contract Address\n\nStream logs from specific contract addresses using `redis_set_filter`:\n\n```json\n{\n \"name\": \"USDC Transfer Monitor\",\n \"description\": \"Monitor USDC transfers on Polygon\",\n \"source\": {\n \"chain\": \"polygon\",\n \"entity\": \"log\"\n },\n \"transforms\": [\n {\n \"type\": \"redis_set_filter\",\n \"set_values\": [\n \"0x3c499c542cef5e3811e1192ce70d8cc03d5c3359\"\n ],\n \"filter_expr\": \"root = this.address\"\n }\n ],\n \"sinks\": [\n {\n \"type\": \"kafka\",\n \"name\": \"usdc-logs\"\n }\n ]\n}\n```\n\n### Example 2: Transform Log Data with JavaScript\n\nParse and enrich log data using a v8 transform:\n\n```json\n{\n \"name\": \"DEX Trade Parser\",\n \"description\": \"Parse DEX swap events and extract trade details\",\n \"source\": {\n \"chain\": \"polygon\",\n \"entity\": \"log\"\n },\n \"transforms\": [\n {\n \"type\": \"redis_set_filter\",\n \"set_values\": [\n \"0xd78ad95fa46c994b6551d0da85fc275fe613ce37657fb8d5e3d130840159d822\"\n ],\n \"filter_expr\": \"root = this.topic0\"\n },\n {\n \"type\": \"v8\",\n \"script\": \"function transform(record) { record.parsed = true; return record; }\"\n }\n ],\n \"sinks\": [\n {\n \"type\": \"kafka\",\n \"name\": \"dex-trades\"\n }\n ]\n}\n```\n\n### Example 3: Monitor ERC-20 Transfers on Base\n\nTrack token transfers on Base using the `erc20_token_transfer` entity:\n\n```json\n{\n \"name\": \"Base Token Transfer Tracker\",\n \"description\": \"Monitor ERC-20 token transfers on Base\",\n \"source\": {\n \"chain\": \"base\",\n \"entity\": \"erc20_token_transfer\"\n },\n \"transforms\": [\n {\n \"type\": \"redis_set_filter\",\n \"set_values\": [\n \"0x833589fcd6edb6e08f4c7c32d4f71b54bda02913\"\n ],\n \"filter_expr\": \"root = this.address\"\n }\n ],\n \"sinks\": [\n {\n \"type\": \"kafka\",\n \"name\": \"base-usdc-transfers\"\n }\n ]\n}\n```\n\n### Example 4: Multiple Contract Addresses\n\nMonitor multiple contracts in a single pipeline:\n\n```json\n{\n \"name\": \"Multi-Token Monitor\",\n \"description\": \"Monitor USDC and USDT on Polygon\",\n \"source\": {\n \"chain\": \"polygon\",\n \"entity\": \"log\"\n },\n \"transforms\": [\n {\n \"type\": \"redis_set_filter\",\n \"set_values\": [\n \"0x3c499c542cef5e3811e1192ce70d8cc03d5c3359\",\n \"0xc2132d05d31c914a87c6611c10748aeb04b58e8f\"\n ],\n \"filter_expr\": \"root = this.address\"\n }\n ],\n \"sinks\": [\n {\n \"type\": \"kafka\",\n \"name\": \"stablecoin-logs\"\n }\n ]\n}\n```\n\n## Tool Reference\n\n### Configuration Management\n\n| Tool | Description |\n| -------------------- | --------------------------------------------------- |\n| `create_beam_config` | Create a new pipeline configuration |\n| `get_beam_config` | Get full details of a config by ID |\n| `list_beam_configs` | List all your pipeline configurations |\n| `update_beam_config` | Update an existing configuration |\n| `delete_beam_config` | Delete a configuration (also tears down deployment) |\n\n### Deployment\n\n| Tool | Description |\n| --------------------------- | ----------------------------------------- |\n| `deploy_beam_pipeline` | Deploy a config to Kubernetes |\n| `teardown_beam_pipeline` | Remove deployed infrastructure |\n| `get_beam_deployment_stats` | Check deployment status and worker health |\n\n### Web Interface\n\nEvery pipeline has a management page at `app.allium.so/build/beam/{config_id}`.\nThe UI URL is returned by `create_beam_config`, `deploy_beam_pipeline`, and `get_beam_deployment_stats`.\n\n## Consuming Your Data\n\nAfter deploying a Beam pipeline with a Kafka sink, the `deploy_beam_pipeline` tool returns everything you need to start consuming:\n\n- **Kafka connection credentials**: bootstrap server, username, and password\n- **Code snippets**: Ready-to-use consumer code in Python (`confluent-kafka`) and TypeScript (`kafkajs`)\n- **Topic names**: Automatically generated as `beam.{config_id}.{sink_name}`\n\nJust copy the provided code snippet, install the relevant Kafka client library, and start receiving data.\n\n## Monitoring\n\n### Checking Deployment Health\n\nAfter deploying, use `get_beam_deployment_stats` to verify:\n\n```json\n{\n \"config_id\": \"abc123\",\n \"status\": \"deployed\",\n \"workers_health\": {\n \"total_workers\": 2,\n \"healthy_workers\": 2,\n \"unhealthy_workers\": 0,\n \"crashing_workers\": 0,\n \"oom_killed_workers\": 0\n }\n}\n```\n\n### Worker Health Fields\n\n| Field | Description |\n| -------------------- | ---------------------------- |\n| `total_workers` | Number of worker pods |\n| `healthy_workers` | Workers running normally |\n| `unhealthy_workers` | Workers not ready |\n| `crashing_workers` | Workers in crash loop |\n| `oom_killed_workers` | Workers killed due to memory |\n\n## Troubleshooting\n\n### Unhealthy Workers\n\nIf `unhealthy_workers > 0`:\n\n1. Check if the pipeline config is valid\n2. Verify source data is available\n3. Update the config and redeploy (no teardown needed)\n\n### OOM Killed Workers\n\nIf `oom_killed_workers > 0`:\n\n1. Simplify transforms to reduce memory usage\n2. Add more aggressive filtering earlier in the pipeline\n3. Contact <support@allium.so> for resource limit adjustments\n\n### Crashing Workers\n\nIf `crashing_workers > 0`:\n\n1. Check v8 script for syntax errors\n2. Verify filter expressions are valid\n3. Review transform logic for runtime errors\n\n## Best Practices\n\n1. **Filter early**: Apply `redis_set_filter` before `v8` transforms to reduce data volume\n2. **Start simple**: Begin with basic filtering, add transforms incrementally\n3. **Monitor after deploy**: Always check `get_beam_deployment_stats` after deployment\n4. **Test transforms**: Validate v8 scripts handle edge cases\n5. **Use descriptive names**: Clear names help track multiple pipelines\n6. **Lowercase addresses/labels/symbols**: Use lowercase when filtering by these values\n7. **Redeploy, don't teardown**: When updating, just redeploy—it's idempotent with zero downtime\n"
}SHA-256: e448472f31b666c82ef06e3261e4d3d7e0424e7bc17372c9ebeeec477da0e993