scieee AI-readable full text Open interactive document viewer

Poster for "On-demand Memory Compression of Stream Aggregates through Reinforcement Learning"

Liu, Jingyu; Gulisano, Vincenzo

Abstract

This is the poster of the paper "On-demand Compression of Stream Aggregates through Reinforcement Learning", which was published at ICPE 2025.

Full text

Engineering'and' Physical'Sciences' Research'Council' Grant'number' EP/X029174/1 Horizon'Europe'2021-2027' Framework'Programme' Grant'Agreement'number'101072456 Disclaimer:+Funded+by+the+European+Union.+ Views+and+opinions+expressed+are+however+those+of+ the+author(s)+only+and+do+not+necessarily+reflect+those+ of+the+EU.+The+EU+cannot+be+held+responsible+for+them. Funded&by the&European&Union Chalmers University of Technology and University of Gothenburg, Sweden ICPE 2025 Research Track Stream Processing and Aggregates Reinforcement Learning Evalua8on Usecases and Setup SPE Controller Environment !state 𝑠 "ac'on π‘Ž RL Agent #reward π‘Ÿ good ac'on posi've reward bad ac'on nega've reward - 𝑖𝑛𝑝𝑒𝑑'π‘Ÿπ‘Žπ‘‘π‘’ - π‘‘β„Žπ‘Ÿπ‘œπ‘’π‘”β„Žπ‘π‘’π‘‘ - π‘œπ‘’π‘‘π‘π‘’π‘‘'π‘Ÿπ‘Žπ‘‘π‘’ - π‘™π‘Žπ‘‘π‘’π‘›π‘π‘¦ - 𝑛/𝑐 π‘Ÿπ‘Žπ‘‘π‘–π‘œ - πΆπ‘ƒπ‘ˆ'π‘π‘œπ‘›π‘ π‘’π‘šπ‘π‘‘π‘–π‘œπ‘› -π‘™π‘Žπ‘‘π‘’π‘›π‘π‘¦ - 𝑛/𝑐'π‘Ÿπ‘Žπ‘‘π‘–π‘œ - #𝑠𝑑𝑒𝑝𝑠'π‘π‘’π‘Ÿ'π‘’π‘π‘–π‘ π‘œπ‘‘π‘’ send data A 8:00 20 A 8:03 15 F fF fF f F fF fF f Input stream Stream Processing Engine (SPE) Can run distributed/in parallel Γ  spread in the Cloud-IoT con'nuum Directed Acyclic Graph Outputs π΄π‘Šπ΄,π‘Šπ‘†, 𝑆, 𝑓 !, 𝑓 "##, 𝑓 $%&, 𝑓 '( Func'on 𝑓 !𝑑 π‘Šπ΄ (window advance) π‘Šπ‘† (window size) Func'on 𝑓 "## Ξ“,𝑑 Func'on 𝑓 $%& Ξ“ Func'on 𝑓 '( Ξ“,𝑑 (opt.) event &me size advance remove 𝒕 output Aggregate Stream 𝑆 ØEnvironment ointerface to connect the SPE and RL Agent ØAgent oimplement training algorithm by Neural Network (DQN) oget an ac'on to interact with the environment ØReward ofeedback to the Agent to reinforce good ac'ons RL agent reward β“΅ β“Ά β“· environment ØLinear Road benchmark oVehicles travelling in highways report their posi'on/speed oEach vehicle reports its posi'on every 30 seconds oAggregate: count the number of non-consecu've stops oWS = 10 mins, WA = 5 secs Ø Synthe;c (stress-test) oData is generated following a sawtooth wave whose peaks’ values and distances are chosen randomly oThe key aYribute is generated from a Gaussian distribu'on with changing πœ‡/𝜎) oAggregate: perform math opera'ons on a random value carried by each tuple oWS = 15 mins, WA = 1 sec ØSetup oJava (OpenJDK 17.0.7), Python 3.7.6 oSPE: Liebre oCompression library: snappy oAgent: openAI Gym o120 episodes with maximum 1000 steps for each On-demand Memory Compression of Stream Aggregates through Reinforcement Learning Comparison discussions for the Agent with different compression levels Scalability discussions for the Aggregate ØWithout an Agent (top): oaverage CPU cons.: 0.33 (Linear Road), 0.59 (Synthetic) oaverage latency: 0.98s (Linear Road), 0.53s (Synthetic) ØWith an Agent (bottom): odiff. in CPU cons. and latency are almost 0 It does not become a scaling bo5leneck for the Aggregate by introducing the Agent. Ø Linear Road (WS = 10 mins, WA = 5 secs) oall baselines are safe except for 𝐷0.0 osimilar policy behaviors except for WEL-OB Ø Synthetic (WS = 15 mins, WA = 1 sec) o𝑛/𝑐 ratio decreases linearly with lower 𝐷 ofine-tune ability (ini'al) state ac'on once per day once per sec send data - πΆπ‘œπ‘šπ‘π‘Ÿπ‘’π‘ π‘  π‘šπ‘œπ‘Ÿπ‘’ ‒𝑒. 𝑔. 𝐷 ↑ 10% -π‘†π‘‘π‘Žπ‘¦ β€’π‘’π‘›π‘β„Žπ‘Žπ‘›π‘”π‘’π‘‘ - πΆπ‘œπ‘šπ‘π‘Ÿπ‘’π‘ π‘  𝑙𝑒𝑠𝑠 ‒𝑒. 𝑔. 𝐷 ↓ 10% oRL-based adaptive memory compression scheme for stream Aggregates oAllowing real-time balancing of performance and memory usage under latency constraints oCapture applicationand data-specific behaviors of Aggregates oHighlight the trade-off between RL training timeliness and policy effectiveness 10:00:00 ac'on 'me β€’If a window hasn’t been updated for a while… β€’Compress it firstly, and later decompress it example About this paper Jingyu Liu, Vincenzo Gulisano infrequently frequently Compress! 'me for the next output 10:05:00 RL Agent SPE ac&on (compress) state, reward Cyclical dependency Ø Agent computes its next ac'on o receive the state and reward first! baseline (X value) baseline (X value) Each 𝐷𝑋#baseline always sets the 𝐷 value to 𝑋 βˆ— π‘Šπ‘†#Each 𝐷𝑋 baseline always sets the 𝐷 value to 𝑋 βˆ— π‘Šπ‘† Output Stream Ø Aggregate waits for the Agent’s ac'on o share a new state and reward first! WEL-OB (WEL -OBlivious ) EL-OB (EL -OBlivious ) L-OB (L-OBlivious) WEL - AW (WEL-AWare) wallclock time (W) βœ˜βœ“ βœ“ βœ“ event time (E) ✘ ✘ βœ“ βœ“ next output (L) (e.g. latency) ✘ ✘ ✘ βœ“ -𝐷𝑄𝑁 policy observa&on Linear Road Synthetic Linear Road Synthe<c (compression threshold) 𝑫= 𝑋 βˆ— π‘Šπ‘†, 𝑋 ∈ 0.0,1.0 : o 𝒍𝒂𝒕𝒆𝒔𝒕𝑻𝑺 βˆ’π’•π’” β‰₯𝑫 Γ  compress (the condition for triggering compression) β€’π‘™π‘Žπ‘‘π‘’π‘ π‘‘π‘‡π‘†: the timestamp (event time) of the latest tuple processed by the Aggregate 𝐴 ‒𝑑𝑠: the timestamp (event time) of the latest tuple that contributed to the window instance β€’e.g., 0.0 βˆ— π‘Šπ‘†: all window instances maintained by the Aggregate 𝐴 are compressed 1.0 βˆ— π‘Šπ‘†: no window instance maintained by the Aggregate 𝐴is compressed tuple window instance Why are four policies? ØThree ways the system makes progress (1) wallclock 'me moves forward (2) event 'me advances (3) get a new state (e.g. latency) measurement β€’(2) implies (1) because event 4me advances only as 𝐴 processes input, which depends on wallclock 4me. β€’(3) implies (2) because latency updates occur only when event 4me advances enough to produce output. Based on these dependencies, four policies can be established: If the state comes from before 10:05:00, the effects on latency are not measured yet… TIME ALIGNMENT (in four policies) up down 100% WS 𝑫 0% WS 100% WS 𝑫 0% WS