influxdata / influxdata/influxdb
The data is deleted after a while, but it must be kept without deleting
- Dominant language
- Rust
- Stars
- 31.7k
- Forks
- 3.7k
- Avg merge
- 13h 37m
- Merged PRs (30d)
- 8
Description
I'm using Influxdb2 with Rust language to save a Binance Futures 24hrTicker stream (!ticker@arr).
I have set the bucket settings to never delete already written data, but it gets deleted after a while. Why is this happening?
Code:
`
use websocket::client::ClientBuilder;
use websocket::OwnedMessage;
use std::time::Instant;
use std::str::FromStr;
use futures::prelude::*;
use influxdb2::models::DataPoint;
use influxdb2::Client;
#[tokio::main]
async fn main() -> Result<(), Box> {
let mut client = ClientBuilder::new("wss://fstream.binance.com/stream")
.unwrap()
.connect(None)
.expect("Failed to connect");
let host = "http://localhost:8086";
let org = "wannafly";
let token = "QPjlMBmeh-fANqZpkC_sScPbW8weLxsw9yl8Q8s7PsUiATo43_plZSathKcRHJ6NOBx9UB3Uh05RixzG4ZT7yA==";
let bucket = "Ticks";
let clientdb = Client::new(host, org, token);
let mut prev_time = Instant::now();
let done = false;
client
.send_message(&OwnedMessage::Text(
r#"{"method": "SUBSCRIBE", "params": ["!ticker@arr"], "id": 1}"#.to_string(),
))
.expect("Failed to send message");
println!("новый цикл: {:?}", prev_time);
for message in client.incoming_messages() {
println!("message: {:?}", message);
println!("внутри цикла");
let message = match message {
Ok(message) => message,
Err(err) => {
println!("Failed to receive message: {:?}", err);
continue;
}
};
match message {
OwnedMessage::Text(text) => {
let tick: serde_json::Value =
serde_json::from_str(&text).expect("Failed to parse JSON");
if let Some(data) = tick.get("data") {
if let Some(tickers) = data.as_array() {
println!("------------------------");
let mut n = 0;
for ticker in tickers {
n = n + 1;
let mut symbol: String;
let mut price: String;
if let Some(s) = ticker.get("s") {
symbol = s.to_string();
symbol = symbol.replace("\"", "");
} else {
symbol = "none".to_string();
}
if let Some(p) = ticker.get("c") {
price = p.to_string();
price = price.replace("\"", "");
match f64::from_str(&price) {
Ok(price) => {
println!("symbol: {}:{:?}", symbol, price);
let points = vec![
DataPoint::builder("Binance")
.tag("Symbol", symbol)
.field("Price", price)
.build()?,
];
clientdb.write(&bucket, stream::iter(points)).await?;
if done {
break;
}
}
Err(err) => {
println!("Price error: {:?}", err);
}
}
} else {
continue;
}
}
let current_time = Instant::now();
let time_diff = current_time.duration_since(prev_time);
prev_time = current_time;
println!("count: {}", n);
println!("Time interval: {:?}", time_diff);
}
}
}
OwnedMessage::Close(_) => {
println!("Connection closed");
break;
}
_ => {}
}
}
println!("конец цикла");
Ok(())
}
`
My sequence of actions:
I run the influxd command
I run the code
Data is being recorded and displayed in the Data Explorer
After a while, the stream is interrupted by println!("Failed to receive message: {:?}", err); This is not an influxdb2 issue. This is a problem with Binance. I just haven't finished the code yet.
I stop Rust. The data is stopped being recorded and stored for some time.
They are missing the next day.
Contributor guide
Research direction
Reproduce the reported sequence using the supplied Rust writer, the influxd command, the configured bucket, and Data Explorer. First verify the bucket's retention settings and determine whether the missing data is caused by storage or by the interrupted Binance stream. Done means identifying the cause of the next-day gap and documenting a reproducible resolution.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- rust
- Domain
- databases
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 25/100