Machine-readable medium for storing a stream data processing program and computer system
Abstract
This record has no abstract on file.
Term
Projected expiry 28 October 2028.
- Priority and filed
- Granted
- Today
- Projected expiry
13 claims: 3 independent, 10 dependent
- 1Stream data input to a computer equipped with a processor and a storage device is received as primary information, and among the received primary information, the primary information for a predetermined period is acquired as a processing target and generated as secondary information. In a stream data processing program that generates time control information indicating that the time has advanced in addition to the primary information, a procedure for accepting the input stream data as primary information and a time when the primary information is received. The time control information generation unit that generates the time information of the above as the time control information obtains the time when the time control information is generated as the next firing time, and the next firing time is set in the storage device as the next firing time holding area. The procedure for storing the time control information in, the procedure for generating the time control information when the current time information becomes the next firing time held in the holding area, and the procedure for receiving the generated time control information. , A stream data processing program characterized by having the computer execute a procedure of acquiring primary information for a predetermined period of the received primary information as a processing target and generating it as secondary information. プロセッサと記憶装置を備えた計算機に入力されたストリームデータを1次情報として受け付けて、前記受け付けた1次情報のうち所定の期間の1次情報を処理対象として取得して2次情報として生成し、時刻が進んだことを表す時刻制御情報を前記1次情報に加えて発生するストリームデータ処理プログラムにおいて、 前記入力されたストリームデータを1次情報として受け付ける手順と、 前記1次情報を受け付けた時点の時刻情報を前記時刻制御情報として生成する時刻制御情報生成部が、前記時刻制御情報を発生する時刻を次回発火時刻として求め、当該次回発火時刻を前記記憶装置に設定された次回発火時刻保持領域に格納する手順と、 現在の時刻情報が前記保持領域に保持された前記次回発火時刻となったときに、前記時刻制御情報を発生する手順と、 前記発生された時刻制御情報を受け付けたときに、前記受け付けた1次情報のうち所定の期間の1次情報を処理対象として取得して2次情報として生成する手順と、 を前記計算機に実行させることを特徴とするストリームデータ処理プログラム。
- 12A processor, a storage device, and an interface are provided, and stream data set in the storage device and input via the interface is acquired as primary information, and the 1 of the acquired primary information for a predetermined period. In a computer system that acquires the following information and generates it as secondary information, and generates time control information indicating that the time has advanced in addition to the primary information, the primary information is acquired and generated as secondary information. The query execution tree analysis unit that extracts the processing module that needs to generate the time control information at a time different from the time when the primary information is received from the query representing the processing content to be performed, the extracted processing module, and the above 1 The time control information generator that generates the time control information calculates the firing time, which is the time when the time control information is generated, from the time information that receives the next information, and sets the firing time in the storage device. The next firing time calculation unit stored in the holding area and a time control information generation unit that generates time information at the time when the primary information is received as the time control information are provided, and the time control information generation unit is provided. A computer system characterized in that the time control information is generated when the current time information reaches the ignition time held in the next ignition time holding area. プロセッサと記憶装置及びインターフェースを備えて、前記記憶装置に設定されて、前記インターフェースを介して入力されたストリームデータを1次情報として取得し、該取得した1次情報のうち所定の期間の前記1次情報を取得して2次情報として生成し、時刻が進んだことを表す時刻制御情報を前記1次情報に加えて生成する計算機システムにおいて、 前記1次情報を取得して2次情報として生成する処理内容を表す問合せから前記1次情報を受け付けた時刻とは異なる時刻に前記時刻制御情報を生成する必用がある処理モジュールを抽出する問合せ実行木解析部と、 前記抽出した処理モジュールおよび前記1次情報を受け付けた時刻情報から前記時刻制御情報を生成する時刻制御情報生成部が前記時刻制御情報を生成する時刻である発火時刻を算出し、該発火時刻を前記記憶装置に設定した次回発火時刻保持領域に格納する次回発火時刻計算部と、 前記1次情報を受け付けた時点の時刻情報を前記時刻制御情報として生成する時刻制御情報生成部と、を備え、 前記時刻制御情報生成部は、 現在の時刻情報が前記次回発火時刻保持領域に保持された前記発火時刻になったときに、前記時刻制御情報を生成することを特徴とする計算機システム。
- 13The stream data generated by the first computer is received as the primary information, the time information generated by the first computer is added to the primary information and transmitted to the second computer, and the second computer is transmitted. Of the primary information received by the computer, the primary information for a predetermined period is acquired and generated as secondary information, and the first computer generates time control information indicating that the time has advanced. In addition, in the computer system that transmits to the second computer, what is the time when the second computer receives the primary information from the query representing the processing content that acquires the primary information and generates it as secondary information? The time control is performed from the query execution tree analysis unit that extracts the processing modules that need to generate the time control information at different times, the processing module extracted by the second computer, and the time information that receives the primary information. Time control for generating information The next firing time calculator that calculates the firing time, which is the time when the information generating unit generates the time control information, and transmits the firing time to the first computer, and the first computer. Receiving the firing time and holding it in the next firing time holding area that holds the firing time, and the next firing time message receiving unit, The time control information generation unit includes a time control information generation unit that transmits the time control information to the second computer at the ignition time when the first computer is held in the next ignition time holding area. A computer system characterized by that. 第一の計算機で生成されたストリームデータを1次情報として受け付けて、前記第一の計算機で生成された時刻情報を前記1次情報に付与して第二の計算機に送信し、 前記第二の計算機が前記受け付けた1次情報のうち所定の期間の1次情報を取得して2次情報として生成し、 前記第一の計算機が、時刻が進んだことを表す時刻制御情報を前記1次情報に加えて第二の計算機に送信する計算機システムにおいて、 前記第二の計算機が前記1次情報を取得して2次情報として生成する処理内容を表す問合せから前記1次情報を受け付けた時刻とは異なる時刻に前記時刻制御情報を生成する必用がある処理モジュールを抽出する問合せ実行木解析部と、 前記第二の計算機が前記抽出した処理モジュールおよび前記1次情報を受け付けた時刻情報から前記時刻制御情報を生成する時刻制御情報生成部が前記時刻制御情報を生成する時刻である発火時刻を算出し、該発火時刻を前記第一の計算機に送信する次回発火時刻計算部と、 前記第一の計算機が該発火時刻を受け付けて、前記発火時刻を保持する次回発火時刻保持領域に保持する次回発火時刻メッセージ受信部と、 前記第一の計算機が前記次回発火時刻保持領域に保持された発火時刻において、前記時刻制御情報生成部が前記時刻制御情報を前記第二の計算機に送信する時刻制御情報生成部と、を備えたことを特徴とする計算機システム。
Independent claims3
214 paragraphs, as filed
The present invention relates to a method of generating time control information inside and between systems in a stream data processing system.
There is an increasing demand for a database management system (hereinafter referred to as DBMS) that executes processing on data stored in a storage device, and a data processing system that processes data that arrives from moment to moment in real time.
A stream data processing system has been proposed as a data processing system suitable for real-time processing of such stream data by defining such data that arrives from moment to moment as stream data. For example, Non-Patent Document 1 discloses a stream data processing system STREAM.
In a stream data processing system, unlike a conventional DBMS, a query (query) is first registered in the system, and the query is continuously executed as data arrives. In STREAM, in order to process stream data efficiently, a concept called a sliding window that cuts out a part of stream data is introduced, and a lifetime is given to the data. A preferable example of a query description language including a sliding window specification is CQL (Continuous Query Language) disclosed in Non-Patent Document 1. CQL is an extension that specifies a sliding window by using parentheses after the stream name in the FROM clause of SQL (Structured Query Language), which is widely used in DBMS.
As for SQL, those disclosed in Non-Patent Document 2 are known. There are two typical methods for specifying the sliding window: (1) the method of specifying the number of data strings to be cut, and (2) the method of specifying the time interval of the data strings to be cut. For example, Rows 50 Preceding shown in Section 2 of Non-Patent Document 1 is a good example of cutting 50 rows of data as a processing target (1), and Range 15 Minutes Preceding is 15 minutes. This is a good example of (2), which cuts out minute data as a processing target. In the case of (1), the survival period of the data is until 50 data arrives, and in the case of (2), it is 15 minutes. The stream data cut by the sliding window is stored in memory and used for query processing.
In stream data processing, in order to realize processing such as event extraction by analysis combining multiple data sources or extraction of events that occurred within a certain period of time, a heart for advancing the time even during the period when no data is generated. It was necessary to periodically generate and process beat tuples (hereinafter abbreviated as HBT (Heart Beat Tuple)) inside the data processing system. The HBT has an HBT flag indicating that it is an HBT, and time information when it occurs.
For example, in a join operation that combines the data sources of two or more inputs, the data sources of two or more inputs are acquired in the order of the oldest time, but the first input is the input from the data source and the second. If there is no input from the data source in the input of, the first input cannot be processed because data before the time of the data source of the first input may be input to the second input. , Waiting for processing will occur. In such a case, the first input can be processed by inputting the HBT having the time information after the time of the data source of the first input to the data source of the second input. , Waiting for processing is eliminated.
As a known method, Non-Patent Document 3 discloses a method in which each query holds two states of yield (data is in the output queue) and more (data is in the input queue), and the operator to be executed next is determined. Has been done. The method disclosed in Non-Patent Document 3 executes an execution tree from the input side to the point where it can be executed, and backtracks to the next executable operator. When backtracking to the input of stream data, Education Time-Stamps (ETS, equivalent to HBT) are played.
Further, as another known method, one logical time is given to a plurality of physical times in a distributed input source, and a method of sorting in a buffer according to a value called Output Bookmark is described in Patent Documents. It is disclosed in 2.<patcit num="1"><text>U.S. Pat. No. 5,495,600</text></patcit><patcit num="2"><text>US2008 / 0072221A1 Event Stream Conditioning</text></patcit><nplcit num="1"><text>By R. Motwani, J. Widom, A. Arasu, B. Babcock, S. Babu, M. Datar, G. Manku, C. Olston, J. Rosenstein, and R. Varma: Query Processing, Resource Management, and Approximation in a Data Stream Management System , In Proc. Of the 2003 Conf. On Innovative Data Systems Research (CIDR), [online], January 2003, [Search October 15, 2008], Internet URL <http //infolab.usc.edu/csci599/Fall2002/paper/DS1_datastreammanagementsystem.pdf></text></nplcit><nplcit num="2"><text>CJ Date, by Hugh Darwen: A Guide to SQL Standard (4th Edition), published by Addison-Wesley Professional, USA, published November 8, 1996, ISBN: 021964260</text></nplcit><nplcit num="3"><text>Optimizing Timestamp Management in Data Stream Management Systems, Yijian Bai, Metal Thakkar, Haixun Wang, Carlo Zaniolo, IEEE 23rd International Conference on Data Engineering, 2007.ICDE 2007,15-20 April 2007 Page: 1334 --1338.</text></nplcit>
<p> For example, in a system for buying and selling stocks, how quickly it can react to fluctuations in stock prices is one of the most important issues of the system, and stock data is temporarily stored in a storage device like a conventional DBMS. Therefore, a method of searching for the stored data cannot respond immediately to the speed of stock price fluctuations, and may miss a business opportunity. For example, Patent Document 1 discloses a mechanism in which a stored query is executed periodically, but for real-time data processing in which it is important to execute a query at the moment when data is input, such as a stock price. Was difficult to apply.</p><p> In Non-Patent Document 2, since HBT is generated and processed periodically inside, the processing timing is constrained by HBT, so that the HBT transmission interval appears as the average latency from data generation to event extraction. In order to reduce the latency, it was necessary to increase the HBT generation rate, which was an overhead for increasing the CPU load and reducing the throughput.</p><p> Further, when performing stream data processing in the second computer using the time information in the first computer in a plurality of computers, it is necessary to transmit the time information in the first computer to the second computer, which is described above. A problem similar to the problem that occurred occurs.</p><p> Even in Non-Patent Document 3, for example, HBT may be generated in addition to data input such as a range window operator and an RStream operator that generate a minus tuple (Negative Tuple method), and the method disclosed in the above document corresponds to this. Can not. There is also a problem that the number of backtracks increases when a large number of queries are registered and the execution tree becomes long.</p><p> Further, in Patent Document 2, since HBT is periodic, it is necessary to increase the HBT generation rate in order to shorten the latency, and there remains a problem that the CPU load increases and the throughput decreases.</p><p> Therefore, the stream data processing system is expected to be applied to applications that require real-time processing, such as financial applications, traffic information systems, distribution systems, traceability systems, sensor monitoring systems, and computer system management.</p><p> In other words, in stream data processing, in order to realize processing such as event extraction by analysis combining multiple data sources or extraction of events that occurred within a certain time, the time is advanced even during the period when no data is generated. It was necessary to periodically generate and process the time control information of the above in the data processing system. However, since the processing timing is determined by the time control information, the time control information transmission interval appears as the average latency (latency within a predetermined time) from the data generation to the event extraction. In other words, if there is executable processing during the period of waiting for the time control information to advance the time, the processing start time in the data processing system is bound by the time control information, and the time to wait for this time control information is the latency. Appears. In order to reduce the latency, it was necessary to increase the generation rate of time control information, which was an overhead due to an increase in CPU load and a decrease in throughput.</p><p> Further, when performing stream data processing in the second computer using the time information in the first computer in a plurality of computers, it is necessary to transmit the time information in the first computer to the second computer, which is described above. A problem similar to the problem that occurred occurs.</p><p> The present invention has been made in view of the above problems, and an object of the present invention is to insert time control information when necessary in stream data processing.</p>
<p> The present invention receives stream data input to a computer, processes the data for which the processing target of the stream data is defined by a window, and adds time control information indicating that the time has advanced to the stream data. In the stream data processing program to be generated, the time information of the received stream data is used as the firing time, which is the time when the time control information generation unit that generates the time control information generates the time control information, in the next firing time holding area. Hold on. Further, a processing module that generates the time control information is extracted from a query representing the processing content of the stream data at a time different from the time when the stream data is received, and the extracted processing module and the stream data are received. The ignition time is calculated from the time information and held in the next ignition time holding area. Then, at the ignition time held in the next ignition time holding area, the time control information generation unit inserts the time control information.</p>
<p> By applying the present invention, it is possible to realize stream data processing with low latency while reducing the amount of time control information.</p>
Hereinafter, an embodiment of the present invention will be described with reference to the accompanying drawings.
FIG. 1 shows the hardware environment of the stream data processing server 100. The stream data processing server 100 includes a CPU 11 that is executed by one computer and performs arithmetic processing, a memory 12 that stores stream data 21 and a program for stream data processing, a disk device 13 that stores data, and a CPU 11 and a disk device. Includes interface 14 to connect 13 and network 112. The stream data processing server 100 may be executed by a plurality of computers.
The data input as the stream data 21 includes sensor nodes such as a temperature sensor node 101 and a humidity sensor node 102, and is connected to the sensor base station 108 and the cradle 109 via the network 105. In addition, the RFID tag 103 is connected to the RFID (Radio Frequency Identification) reader 110 via the network 106. Further, the mobile phone 104 is connected to the mobile phone base station 111 via the network 107. Further, the stock providing server 118 that distributes the stock price information may be input as the stream data 21.
The network 112 includes a relay computer 113 and a stream that execute applications such as the sensor base station 108, the cradle 109, the RFID reader 110, the mobile phone base station 111, the stock providing server 118, the sensor net middleware, and the RFID middleware. The computer 115 that inputs a command to the data processing server 100 and the computer 117 that uses the output result 23 output by the stream data processing server 100 are connected.
Sensor base station 108, the temperature sensor node 101, the measurement result of the humidity sensor nodes (e.g., temperature and humidity) exiting the by force, RFID reader 110 outputs the information of the RFID tag 103 to read. The mobile phone base station 111 outputs the information of the mobile phone 104. These outputs are input to the stream data processing server 100 as stream data 21. Even if stream data 21 is directly input to the stream data processing server 100 via the network 112 of the sensor base station 108, the cradle 109, the RFID reader 110, the mobile phone base station 111, and the stock information providing server 118. , Stream data 21 may be input to the stream data processing server 100 after processing by the relay computer 113.
Further, the command 22 commanded by the user 114 or the command 22 generated by the computer 115 is input to the stream data processing server 100 via the network 112. The output result 23, which is the result processed by the stream data processing server 100, is output to the computer 117 used by the user 116 via the network 112.
Here, the stream data processing server 100, the relay computer 113, the computer 115, and the computer 117 are configured by any computer system such as a personal computer and a workstation, and may be the same computer or different computers. Absent. Further, the user 114 and the user 116 may be the same user or different users. In addition, networks 105, 106, 107, and 112 are connected via Ethernet (registered trademark), optical fiber, FDDI (Fiber Distributed Data Interface), wireless, etc., local area network (LAN), or the Internet, which is slower than LAN. It can be a included wide area network (WAN), a public telephone network, or a similar technology that will be invented in the future.
Here, the storage device 15 shown in FIG. 1 is composed of a predetermined area of the memory 12 and a predetermined area of the disk device 13. The stream data 21 is mainly stored in the storage device 15 on the memory 12, and enables high-speed retrieval in response to an inquiry. For the stream data 21 that changes from moment to moment, the data to be searched is stored in the storage device 15 on the memory 12, and the data that has been searched is stored in the storage device 15 on the disk device 13. Can be done. For example, if the stream data 21 is the measured value of the temperature sensor node 101 (temperature, etc.), the measured value that the user 114 wants to monitor is today's value, and there is no problem even if the values before yesterday cannot be searched at high speed. .. Therefore, today's measured values can be stored in the memory 12, and the measured values before yesterday can be stored in the disk device 13 as an archive. Here, the storage device 15 may be any storage medium such as a memory, a disk, a tape, or a flash memory. Further, the storage device 15 may have a hierarchical structure composed of a plurality of storage media. Further, the storage device 15 may be a similar technique invented in the future.
In FIG. 1, the stream data processing server 100 inputs information transmitted in real time from a sensor base station 105, an RFID reader 110, or an application executed on the relay computer 113 via the I / F 14 as stream data 21. Then, based on the command 22 entered by the user 114 or the application running on the computer 115, the input stream data 21 is converted into meaningful information to generate the output result 23, which is the user 116 or the computer. 117 A computer (or server) that executes stream data processing provided to applications running on it. Stream data is multiple stream data 21<sub>1</sub>、21<sub>2</sub>、・・・、21<sub>n</sub>Consists of.
The calculator 115 is connected to the stream data processing server 100 via the network 112. Here, the application executed on the computer 115, the application executed on the computer 116, and the application executed on the computer 117 may be the same application or different applications.
Here, the stream data 21 handled in the present embodiment is different from the stream used in the distribution of video and audio, and one stream data corresponds to significant information. Further, the stream data 21 received by the stream data processing server 100 from an application executed by the stream data processing server 100 on the sensor base station 105, the RFID reader 110, or the relay computer 113 is continuous or intermittent. , Different product information and different elements are included for each stream data.
FIG. 2 is a block diagram showing a stream data processing server to which one embodiment of the present invention is applied and a related system configuration.
The stream data processing server 100 is composed of a CPU 11, a memory 12, a DISK 13, and an I / F 14, and the memory 12 is composed of an operating system (OS) 200, a command input unit 210, and a stream data processing unit 220. The stream data processing unit 220 and the command input unit 210 are composed of programs and are stored in a storage medium such as DISK 13. When performing stream data processing, the CPU 11 loads the stream data processing unit 220 and the command input unit 210 into the memory 12 and executes them.
The command input unit 210 receives a command input by the user 114 or an application executed on the computer 115. Next, the stream data processing unit 220 converts the information of the stream data 21 into significant information based on the query representing the processing content for converting the stream data received by the command input unit 210 into significant information. Output from.
The outline of the present invention will be described with reference to FIG. The stream data processing server 100 reads the stream data 21 based on the query input by the user 114 or the application executed on the computer 115, the query execution unit 230 reads the stream data 21, converts it into meaningful information, and then outputs the result 23. Is output. Here, the significant information is, for example, the measured value of the temperature sensor node 101 shown in FIG. 1 as the average value of the temperature for a certain period of time because the users 114 and 116 cannot understand the measured value data series as it is. This is the converted information.
Hereinafter, the configuration of the stream data processing server 100 will be described in detail.
The command input unit 210 includes an interface (hereinafter, I / F) that receives a command 22 commanded by the user 114 from the computer 115 or a command 22 input from an application executed on the computer 115. When the command is a command related to stream data, the stream data processing unit 220 receives a stream data registration / change command representing a command for registering or changing the stream data input from the command input unit 210, and receives the stream data. Update the table that manages (not shown).
When the command is a command related to a query, the stream data processing unit 220 manages the query when it receives a query registration / change command representing a command for registering or changing the query input from the command input unit 210. The table (not shown) is updated to generate or change the query execution unit 226 that represents the processing content of the stream data corresponding to the query command. Further, the stream data processing unit 220 transmits the generated execution tree 226 to the inquiry processing unit 226 and stores it.
The stream data processing unit 220 includes a stream data receiving unit 221, a query execution tree scheduler 222, a next firing time calculation unit 223, an HBT (Heart Beat Tuple) generation unit 224, a query execution tree analysis unit 225, a query execution unit 226, and an input stream. It consists of data retention buffer 231, system time retention area 232, data processing time retention area for HBT generation 233, next firing time retention area 234, query execution tree analysis result management table 235, operator concatenation queue 236, and output result retention buffer 237. To.
The input stream data holding buffer 231 is a buffer that holds the stream data 21 input to the stream data processing server 100 via the I / F 14. The system time holding area 232 is an area for holding the current time of the system. In the present embodiment, the current time of the system is the absolute time information possessed by the stream data processing server 100 (for example, the current time managed by the OS 200). Further, the current time of the system may be a value updated by time information input from the outside of the stream data processing server 100. For example, the time information given to the stream data 21 may be used to update the current time of the system when the stream data 21 including the latest time information is input.
The stream data receiving unit 221 acquires the data of the input stream data holding buffer 231, adds the current time of the system held in the system time holding area 232 to the stream data 21, and outputs the data to the HBT generating unit 224. ..
The query execution unit 226 represents the contents of processing the stream data 21, and has a tree structure of processing modules such as window operation, selection operation, projection operation, join operation, and aggregation operation. Hereinafter, the processing module is referred to as an operator, and the tree structure is referred to as an execution tree. The execution tree in the query execution unit 226 is generated when a command related to the query is input to the command input unit 210. The query execution unit 226 receives the data output from the HBT generation unit 224, and stores the result processed by each operator of the execution tree in the output result holding buffer 237. The intermediate result processed by each operator is stored in the operator concatenation queue 236. The execution tree of the query execution unit 226 is the same as that disclosed in Japanese Patent Application Laid-Open No. 2008-123426 proposed by the applicant of the present application, and is not described in detail here.
The query execution tree scheduler 222 controls the execution order of the operators in the HBT generation unit 224 and the query execution unit 226.
The operator concatenation queue 236 is a buffer that holds the intermediate result processed by each operator. The output result holding buffer 237 is a buffer for storing the processing result output by the query execution unit 226. The output result stored in the output result holding buffer 237 is output to the computer 117 shown in FIG. 1 via the I / F14.
The query execution tree analysis unit 225 analyzes the execution tree in the query execution unit 226, extracts an operator who generates the HBT at a time different from the time when the stream data is received, and sets the operator and the query. The information is stored in the query execution tree analysis result management table 235. The query execution tree analysis result management table 235 is a table that stores the results analyzed by the query execution tree analysis unit 225.
The next firing time calculation unit 223 refers to the query execution tree analysis result management table 235, and is based on the time information of the input stream data and the setting information of the query stored in the query execution tree analysis result management table 235. The next firing time is calculated, and the next firing time is held in the next firing time holding area 234.
Here, the next firing time is a time when the inquiry execution unit 226 starts processing other than the arrival of the stream data 21.
The next firing time calculation unit 223 is called when the query execution unit 226 executes the operator extracted by the query execution tree analysis unit 225. The next firing time holding area 234 is an area for holding the next firing time calculated by the next firing time calculation unit 223.
The HBT generation unit 224 acquires the stream data 21 from the stream data reception unit 221, stores the reception time of the stream data 21 in the HBT generation data processing time holding area 233, and outputs the stream data 21 to the query execution unit 226. .. The HBT generation data processing time holding area 233 is an area for holding the final time when the HBT generation unit 224 processed the stream data 21.
Further, the HBT generation unit 224 refers to the system time holding area 232, the HBT generation data processing time holding area 233, and the next firing time holding area 234, and the details will be described later in the system time holding area 232. The firing time starts from the current time of the system to be held, the next firing time held in the next firing time holding area 234, and the last time processed by the HBT generation unit held in the HBT generation data processing time holding area 233. The HBT is output to the query execution unit 226.
Here, the stream data 21, the output result 23, the HBT described later, and the temporarily stored data held by the operator for processing may be any data format such as a tuple format (record format), an XML format, or a CSV file. An example of using the tuple format will be described below.
Here, the stream data 21, the output result 23, the HBT described later, and the temporarily stored data held by the operator for processing do not need to have a data entity, and a part or all of the data is a data entity. It may include a pointer to point to.
3a and 3b are diagrams schematically showing an example of a suitable data format of the stream data 21. In the illustrated example, the temperature stream data (S1) 21 output by the temperature sensor node 101 in FIG. 3a, respectively.<sub>1</sub>, Humidity stream data (S2) 21 output by the humidity sensor node 102 in Fig. 3b.<sub>2</sub>Is shown.
Temperature stream data (S1) 21 in Figure 3a<sub>1</sub>Is a record format, and the temperature sensor ID column 302, the device ID column 303, the temperature column 304, and the system time stamp column 305 constituting the record correspond to segments, and the temperature sensor ID column 302, the device ID column 303, and the above The combination of the temperature column 304 and the system time stamp column 305 is referred to as a tuple 301.
Here, the value of the system time stamp column 305 is the temperature stream data (S1) 21.<sub>1</sub>Is input to the stream data processing server 100, and the time information when the stream data processing server 100 arrives at the stream data processing server 100 is given.
The system time stamp column 305 may use the time information given before being input to the stream data processing server 100. For example, the value of the system time stamp column 305 is the temperature stream data (S1) 21.<sub>1</sub>Is the time information given before being input to the stream data processing server 100, and is given by the temperature sensor 101, the sensor base station 106, or the relay computer 113 on an application such as a sensor net middleware. May be done.
Humidity stream data (S2) 21 in Figure 3b<sub>2</sub>Is a record format, and the humidity sensor ID column 312, the device ID column 313, the humidity column 314, and the system time stamp column 315 that make up the record correspond to segments, and the humidity sensor ID column 312, the device ID column 313, and the above The combination of the humidity column 314 and the system time stamp column 315 is a tuple 311.
Here, the value of the system time stamp column 315 is the humidity stream data (S2) 21.<sub>2</sub>Is input to the stream data processing server 100, and the time information when the stream data processing server 100 arrives at the stream data processing server 100 is given.
The system time stamp column 315 may use the time information given before being input to the stream data processing server 100. For example, the value of the system time stamp column 315 is the humidity stream data (S2) 21.<sub>2</sub>Is the time information given before being input to the stream data processing server 100, and is given by the temperature sensor 101, the sensor base station 106, or the relay computer 113 on an application such as a sensor net middleware. May be done.
FIG. 4 is a description example of a suitable command when registering / setting the stream data 21 in the stream data processing server 100 in the command input unit 210.
The stream registration command 401 is registered from the application 116 instructed by the user 114 from the computer 115 or executed on the client computer 115 through the command input unit 210. The stream registration command 401 is the temperature stream data (S1) 21.<sub>1</sub>Registers stream data consisting of a temperature sensor ID that holds an integer type (int type), a device ID that holds an integer type (int type), and a temperature that holds a floating point type (double type). Indicates that it is a command. These correspond to the temperature sensor ID column 302, the device ID column 303, and the temperature column 304 shown in FIG. 3, respectively.
Further, the stream registration command 402 uses the humidity stream data (S2) 21.<sub>2</sub>Registers stream data consisting of a humidity sensor ID that holds an integer type (int type), a device ID that holds an integer type (int type), and humidity that holds a floating point type (double type). Indicates that it is a command. These correspond to the humidity sensor ID column 312, the device ID column 313, and the temperature column 314 shown in FIG. 3, respectively.
In the command input unit 210, the command for registering and setting the stream data 21 in the stream data processing server 100 may be transformed into a table format for managing the command and stored in the storage medium.
In the present embodiment, the system time stamp column 305 and the system time stamp column 315 are automatically included, but in the stream registration command 411, "register stream temperature stream (time stamp timestamp, temperature sensor ID)". It may be explicitly specified as "int, device ID int, temperature double);".
Further, in the present embodiment, an example in which a command is registered in a command line interface (CLI) format is shown, but the present invention is not limited to this. For example, the graphic user interface (GUI) may be used to input the same meaning as above, or the input may be in a tabular format, a setting file, or an XML file. The same applies to the following commands.
Further, in the present embodiment, the time stamp is expressed in the form of hours and minutes such as "9:00", but the date and minutes such as "2007/9/21 9:00:00 JST" and Other formats typified by the format including seconds may be used. The same shall apply in the following figures.
FIG. 5 is a description example of a suitable command when registering and setting an inquiry registration command in the stream data processing server 100 in the command input unit 210.
The inquiry registration command 501 is registered from the user 114 or the application 116 executed on the client computer 115 through the command input unit 210.
The inquiry registration command 501 is the temperature stream (S1) 21.<sub>1</sub>Last 2 minutes ([Range 2 minute]), and said humidity stream (S2) 21<sub>2</sub>For the latest one for each humidity sensor ID ([Partition by S1 temperature sensor ID rows 1]), the temperature stream (S1) 21<sub>1</sub>The condition that the temperature of is 20 degrees or more (S1. Temperature> = 20) and the humidity stream (S2) 21<sub>2</sub>Satisfies the condition that the humidity of is 60% or more (S2. Humidity> = 60), and the temperature stream (S1) 21<sub>1</sub>Temperature sensor ID and the humidity stream (S2) 21<sub>2</sub>If the humidity sensor ID of the above matches, the temperature stream (S1) 21<sub>1</sub>Tuple and said humidity stream (S2) 21<sub>2</sub>(WHERE S1. Temperature sensor ID = S2. Humidity sensor ID), for each device ID (GROUP BY S1. Device ID) Average temperature (Avg (S1. Temperature)), and average humidity (Avg (humidity)) is calculated, only the increment of the tuple of the temperature sensor ID, the average temperature, and the average humidity is streamed (ISTREAM), delayed by 1 minute (<1 minute>), and output. Indicates that the query represents.
The inquiry registration command 502 indicates that the "[Range 2 minute]" of the inquiry registration command 501 is changed to "[Jumping 10 minute]" and the processing target is switched at 10-minute intervals. For example, data entered at 9:01 in the time window (Range 2 minute) will be processed until 9:03, and data entered at 9:04 will be processed until 9:06, while jumping windows. For (Junping 10 minute), both the data input at 9:01 and the data input at 9:04 will be processed at 9: 00-9: 10 (not including 9:10), and at 9:10 9: 10-9: 20 will be processed. In addition, the query registration command 502 changes the "ISTREAM () <1 minute>" of the query registration command 501 to "RSTREAM [5 minute]" and outputs a tuple set of average values every 5 minutes. To do.
The command for registering and setting the query registration command in the stream data processing server 100 in the command input unit 210 may be transformed into a table format for managing the command and stored in the storage medium.
FIG. 6 is an explanatory diagram showing an example of the query execution unit 226.
The query execution unit 226 represents the query execution unit 226 generated when the query registration command 501 shown in FIG. 5 is executed. The query execution unit 226 is composed of an operator that performs processing and an operator connection queue 236 that connects the operators. In this explanatory diagram, the left end is the input and the right end is the output. Stream data 21 is input as input data. Here, among the data strings input as the stream data, each data is referred to as a tuple 601. The tuple 601 is processed by the operator and stored in the operator concatenation queue 236. The tuple 601 is processed by the stream data receiving unit 221 and the HBT generating unit 224 shown in FIG. 2, and then input to the query executing unit 226 via the operator concatenation queue 236. In addition, the processing result of the query of the execution tree 241 is output as the output result 23. It is also possible to re-input the output result 23 as another stream data 21.
The type of the operator differs depending on the processing content. The window operators 611 and 612 shown in FIG. 6 specify the number of data strings from the stream data 21 or specify the time interval of the data strings to be cut, cut the data strings, and convert the stream data into a tapple set. Perform processing. The cut tuples are held in window operators 611, 612. The selection operators 613 and 614 perform a process of determining whether or not to output the tuples 301 and 311 shown in FIG. 3 based on preset conditions. The join operator 615 performs a process of joining stream data 21 having two or more inputs under a predetermined condition. The aggregation operator 616 performs preset aggregation processing such as total, average, maximum, and minimum. The streaming operator 617 performs a process of converting the tuple set into stream data 21 as the output result 23. In addition to the operators shown in FIG. 6, there are projection operators and the like that perform processing to output only a part of the columns of the tuple 203.
Execution tree 241<sub>1</sub>Represents the execution tree 241 generated when the query registration commands 501 and 502 shown in FIG. 5 are executed. Execution tree 241<sub>1</sub>Is the temperature stream data 21<sub>1</sub>, Humidity stream data 21<sub>2</sub>Is input. The window operator 611 has temperature stream data 21<sub>1</sub>The last 2 minutes ([Range 2 minute]) of is held in the window operator 611, and the tuple newly entered in the window and the tuple exited from the window are output to the selection operator 613.
The window operator 612 has humidity stream data 21<sub>2</sub>The latest one for each humidity sensor ID ([Partition by S1 humidity sensor ID rows 1]) is held in the window operator 612, and the tuple newly entered in the window and the tuple exited from the window are output to the selected operator 614. To do.
The selection operator 613 outputs a tuple satisfying the condition that the temperature is 20 degrees or more (S1. Temperature> = 20) to the tuple operator 615 with respect to the tuple output by the window operator 611.
The selection operator 614 outputs a tuple satisfying the condition (S2. Humidity> = 60) that the humidity is 60% or more with respect to the tuple output by the window operator 612 to the coupling operator 615.
When the temperature sensor ID of the tuple output by the selection operator 613 and the humidity sensor ID of the tuple output by the selection operator 614 match, the coupling operator 615 combines the tuples (WHERE S1. Temperature sensor). ID = S2. Humidity sensor ID), output to the tabulation operator 616. The join operator 616 holds the tuples output by the selection operator 613 and the selection operator 614 in the temporary storage area in order to select the tuple to be combined. The tuple held in the temporary storage area may be an entity of data or data including pointers to window operators 611 and 612.
The aggregation operator 616 has the temperature stream data 21 for the output tuple of the join operator 615.<sub>1</sub>Calculate the average temperature (Avg (S1. Temperature)) and the average humidity (Avg (S2. Humidity)) for each device ID of (GROUP BY S1. Device ID), and calculate the temperature sensor ID and average temperature. The average value of value and humidity is output to the streaming operator 617. The aggregation operator 616 holds a tuple for calculating the aggregation value in the temporary storage area.
The streaming operator 617 streams only the increment of the output tuple of the aggregation operator 616 and delays it by 1 minute (ISTREAM <1 minute>), and then outputs the result 23.<sub>1</sub>Output as.
Execution tree 241 above<sub>1</sub>An example of the operation of is as follows.
For example, the stream data 21<sub>1</sub>It is assumed that the tuple (temperature sensor ID, device ID, temperature) = (1001, 201, 23 ° C) is input at 9:00. In the stream data receiving unit 221, the current system time information is added to the system time stamp column 305 shown in FIG. 3 for the tuple, and the tuple 601<sub>1</sub>(Temperature sensor ID, device ID, temperature, system time stamp) = (1001, 201, 23 ° C, 9:00) is output.
In addition, HBT generator 224<sub>1</sub>Is a heartbeat tuple (HBT) 604 for advancing the time in the inquiry execution unit 226 even during the period when no data is generated.<sub>1</sub>Is output. HBT604<sub>1</sub>Has an HBT flag indicating that it is an HBT, and a system time stamp, for example, HBT604 output at 9:03.<sub>1</sub>Is "(HBT, 9:03)". When the HBT is received by each operator, the operator processing time managed by each operator is updated and stored in the next operator concatenation queue 236. The method of generating HBT will be described later.
Tuple 601<sub>1</sub>In the window operator 611, the period (survival period) to be processed is from 9:00 to 9:02 in the window cut in the past 2 minutes ([Range 2 minute]).
Window operator 611 is a tuple 601<sub>1</sub>602 with a plus flag indicating that the tuple's survival has started at the beginning of the tuple's survival to represent the survival of the tuple.<sub>1</sub>603 with a minus flag indicating that the tuple's life has expired at the end of the tuple's life<sub>1</sub>Is output. The tuple plus tuple 602<sub>1</sub>Is "(+, 1001, 201, 23 ° C, 9:00)", minus tuple 603<sub>1</sub>Is "(-, 1001, 201, 23 ° C, 9:00)".
The window operator 611 is a plastic tuple 602.<sub>1</sub>Is output when it is received, and the output tuple is retained in the temporary storage area. And the HBT604<sub>1</sub>When you receive HBT604<sub>2</sub>Is stored in the operator concatenation queue 236. Further, the operator processing time managed by the window operator 611 is updated from "9:00" to "9:03", and the minus tuple 603 is used.<sub>1</sub>Can be output. HBT604 output at 9:03<sub>1</sub>Said minus tuple 603 using<sub>1</sub>When outputting, the time from 9:02 to 9:03 is the inquiry execution unit 226<sub>1</sub>It appears as the latency of. This latency is the latency described in the above-mentioned task. The minus tuple 603<sub>1</sub>In order to reduce the latency to output, the window operator 611 needs to receive the HBT at 9:02.
Minus tuple 603<sub>1</sub>Is sometimes referred to as the Negative Tuple. In addition to this method, any method may be used to express the survival period, such as embedding the end of the survival period in the tuple for processing.
In addition, the HBT604<sub>1</sub>Is also required for operators that handle two or more inputs, such as the join operator 615. For example, join operator 615 has operator join queue 236 in chronological order.<sub>2</sub>, Operator concatenation queue 236<sub>3</sub>Get tuples from.
For example, Plastaple 602<sub>2</sub>"(+, 1001, 201, 21 ° C, 8:40)" is the operator connection queue 236<sub>2</sub>Medium, Plastaple 602<sub>3</sub>"(+, 1001, 201, 69%, 8:30)" is the operator connection queue 236<sub>3</sub>If inside, said Plastuple 602<sub>3</sub>First, the operator connection queue 236<sub>3</sub>The join operator 615 gets it from the inside. Next, the join operator 615 transfers the plastic tuple 602.<sub>2</sub>The operator connection queue 236<sub>2</sub>I try to get it from the inside, but the operator connection queue 236<sub>3</sub>Since data before 8:40 may be stored in the above-mentioned Plastuple 602<sub>2</sub>Cannot be obtained. Here, HBT604<sub>3</sub>"(HBT, 8:45)" is the operator connection queue 236<sub>3</sub>When stored in, the join operator 615 is the plastic 602.<sub>2</sub>Is acquired and the process is executed.
In addition, the aggregation operator 616 may generate tuples called ghosts, which have the same time stamps for the start and end of survival and have no survival period. For example, when the value of the humidity column of the tuple to be aggregated is "64, 66, 68, 70, 72", the average value of the humidity column is "68". Here, at the time stamp "8:20", minus tuple 603<sub>2</sub>When "(-, 64, 8:20)" arrives, the average humidity is "69" because it is the average of "66, 68, 70, 72", and it is a plaster with an aggregated value of "69". Is generated. However, the same time stamped Plus Tuple 602<sub>4</sub>When "(+, 74, 8:20)" arrives, the average humidity is "66, 68, 70, 72, 74", so it is "70" and has a total value of "69". A minus tuple and a plus tuple with an aggregate value of "70" are generated. That is, a tuple having the aggregated value "69" is a tuple having no survival period. When removing ghosts with the aggregation operator 611, a tuple with a system time stamp after 8:20 may arrive, for example, HBT604 with a system time stamp of 8:25.<sub>3</sub>When "(HBT, 8:45)" arrives, the aggregated value of 8:40 is fixed, so the plaster with the aggregated value "70" can be output.
In the present embodiment, the ghost is removed by the aggregation operator 616, but it may be removed by another operator. For example, the streaming operator 617 may remove the ghost. In addition, all aggregate operators 617 have a ghost removal function, but by inputting with command line interface (CLI) format, graphic user interface (GUI), tabular format, configuration file, XML file, ghost removal function for each operator May be turned on and off. Moreover, the average having a ghost removal function may be described in the query such as AVG_G.
In addition, the plastic tuple and HBT in the execution tree shown in this embodiment do not reverse the time in the operator concatenated queue and on the operator having a parent-child relationship in the graph structure, and the time stamp on the input side is new and output. The time stamp on the side becomes the old value.
FIG. 7 is an explanatory diagram showing a configuration example of the query execution tree analysis result management table 235.
The target operator column 701 stores the operator extracted by the query execution tree analysis unit 225 of FIG. The setting item column 702 stores setting items such as the window size, delay size, and output interval set for the operator extracted by the query execution tree analysis unit 225. In the last execution time column 703, in the case of the jumping window or the streaming operator (RStream) that outputs at regular intervals, the operator executes the previous process and stores the output time.
For example, line 704 represents the query execution tree analysis result management table 235 of the operator extracted as a result of analyzing the execution tree generated from the query registration command 501 shown in FIG. The details of the extraction method will be described later.
In line 704, the value of the target operator column 701 is "window operator 611 (Range Window)", the value of the setting item column 702 is "sliding window size = 2 minutes", and the value of the previous execution time column 703 is "-". It represents that.
Here, the table for managing the query execution tree analysis result may be in any format such as a configuration file and an XML file, in addition to the table format shown in FIG. The same applies to the following tables.
FIG. 8 is a flowchart showing the entire processing of the stream data processing server 100. This process is started when the administrator of the stream data processing server 100 starts the stream data processing unit 220 or the like.
First, the query execution tree analysis unit 225 shown in FIG. 2 extracts a firing operator from the query execution tree of the registered query and registers it in the query execution tree analysis result management table 235 (802). The details of the process of step 802 will be described later in FIG. The ignition operator refers to an operator who needs to start processing without waiting for the arrival of stream data 21, and is an operator having a time constraint related to processing, such as an operator who outputs processing results at predetermined time intervals. is there.
Next, the time information for which the stream data 21 is received by the HBT generation unit 224 is held in another HBT generation data processing time holding area 233 (803). The details of the process of step 803 will be described later with reference to FIG.
Next, the next firing time calculation unit 223 calculates the next firing time and holds it in the next firing time holding area 234 (804). The details of the process of step 804 will be described later with reference to FIG.
Next, the HBT is inserted (generated) at the ignition time held in the next ignition time holding region 234 by the HBT generation unit 224 (805). The details of the process of step 805 will be described later with reference to FIG.
Next, it is determined whether or not the system termination command has been accepted by the command input unit 210 (806). If NO is determined in step 806, the process returns to step 802, and if YES is determined in step 806, the processing of the stream data processing server 100 is terminated (807).
FIG. 9 is a flowchart showing the extraction and registration process of the ignition operator in step 802 shown in FIG.
The query execution tree analysis unit 225 repeats the processes from step 903 to step 912 shown below for all the operators in the query execution unit 226 shown in FIG. 2 (902).
First, it is determined whether or not the target operator is a sliding window operator (Range Window) representing time (903). If YES is determined in step 903, the target operator and the sliding window size are registered in the query execution tree graph analysis result management table 235 (904).
When the step 904 is completed, or when the determination is NO in the step 903, it is determined whether or not the target operator is a jumping window operator representing time (905). If YES is determined in step 905, the target operator and the jumping window size are registered in the query execution tree graph analysis result management table 235 (906).
When the step 906 is completed, or when the determination is NO in the step 905, it is determined whether or not the target operator is a streaming operator (IStream, DStream, IDStream) that causes a delay representing time (907). If YES is determined in step 907, the target operator and the delay size are registered in the query execution tree graph analysis result management table 235 (908).
When the step 908 is completed, or when the determination is NO in the step 907, it is determined whether or not the target operator is a streaming operator (RStream) that outputs at regular intervals representing time (909). If YES is determined in step 909, the target operator and the output interval are registered in the query execution tree graph analysis result management table 235 (910).
When the step 910 is completed, or when the determination is NO in the step 909, the target operator has an operator (Sum, Count, Average, Min, Max, Median, Variable, Standard Deviation, Limit) having a function of removing the ghost. ) Whether or not (911). If YES is determined in step 911, the query execution tree graph analysis result management table 235 has the target operator and a serial number under the minimum time unit (for example, 1 millisecond, 1 nanosecond, 1 millisecond). Etc.) are registered (912).
When the step 912 is completed, or when NO is determined in the step 911, the process returns to the step 902, and the steps 903 to 912 are repeated again. When the processing for all operators is completed, the processing in step 802 is completed (913).
In the following, the query registration command 501 shown in FIG. 5, the query registration command 502, and the query execution unit 226 shown in FIG. 6 are shown below.<sub>1</sub>An example is shown in which the query execution tree analysis result management table 235 shown in FIG. 7 is created by the process of FIG. 9 above.
When the query registration command 501 shown in Fig. 5 is registered, the query execution unit 226 shown in Fig. 6 is registered.<sub>1</sub>Is generated. The stream data processing unit 220 is the query execution unit 226.<sub>1</sub>For this, the flowchart of FIG. 9 is executed.
First, since the window operator 611 is a sliding window operator (Range window, [Range 2 minute]) representing time, it is determined to be YES in step 903, and "window operator 611" and "sliding window size = 2 minutes" are determined. Is registered in row 704 of the query execution tree analysis result management table 235.
Next, since the window operator 612 is a group-by-group window (Partitioned window, [Partition by S1 humidity sensor ID rows 1]), all of the steps 903, 905, 907, 909, and 911 are determined to be NO. .. Similarly, the selection operator 613, the selection operator 614, and the join operator 615 are all determined to be NO in steps 903, 905, 907, 909, and 911.
Since the aggregation operator 616 is an operator (Avg (S1. Temperature), Avg (S2. Humidity)) having a function of removing ghosts, it is determined as YES in step 911, and the "aggregation operator 616" and "time" are determined. "Minimum unit" is registered in row 705 of query execution tree analysis result management table 235.
Since the streaming operator 617 is a streaming operator (ISTREAM () <1 minute>) that generates a delay, it is determined as YES in step 907, and "streaming operator 617" and "delay size = 1 minute". Is registered in row 706 of the query execution tree analysis result management table 235.
If the query registration command 502 shown in FIG. 5 is registered, the query execution unit 226 shown in FIG. 6 is registered.<sub>1</sub>The same query execution unit 226 as in is generated, and the flowchart of FIG. 9 is executed.
Since the description of [Jumping 10 minute] of the query registration command 502 generates a jumping window operator (Jumping window) representing the time, it is determined as YES in step 905, and "window operator xxx" and "jumping window size" are generated. = 10 minutes is registered in row 707 of the query execution tree analysis result management table 235.
Further, since the description of RSTREAM [5 minute] of the query registration command 502 generates a jumping window operator (RStream) representing the time, it is determined as YES in step 909, and "window operator yyy" and "output interval" are generated. = 5 minutes is registered in row 708 of the query execution tree analysis result management table 235.
In this embodiment, it is assumed that all aggregation operators have a ghost removal function, but each operator can be input by command line interface (CLI) format, graphic user interface (GUI), tabular format, configuration file, and XML file. The ghost removal function may be turned on and off. Moreover, the average having a ghost removal function may be described in the query such as AVG_G.
By the above processing, the firing operator who has a time constraint to start the processing is registered in the query execution tree graph analysis result management table 235.
FIG. 10 is a flowchart showing the stream data input and registration process of step 803 shown in FIG.
First, the HBT generation unit 224 determines whether or not the stream data reception unit 221 shown in FIG. 2 has received the stream data 21 (1002). If YES is determined in step 1002, the input time information of the received stream data 21 is registered in the next firing time holding area 234 other than the HBT generation unit 224 that received the stream data 21 (1003), and HBT is generated. Part 224 updates the value of the input time information of the stream data 21 that received the final data input time of the data processing time holding area 233 for HBT generation (1004).
When the step 1004 is completed, or when NO is determined in the step 1002, the process of the step 803 is terminated (1005).
FIG. 11 is a flowchart showing the next firing time calculation process of step 804 shown in FIG.
First, the next firing time calculation unit 223 determines whether or not the target operator is a sliding window operator (Range Window) representing time (1102). If YES is determined in step 1102, the time information given to the tuple input to the target operator in the query execution unit 226 in the next firing time holding area 234 + the query execution tree graph analysis result management. Register the value of the sliding window size registered in Table 235 (1103).
When the step 1103 is completed, or when NO is determined in the step 1102, the next firing time calculation unit 223 determines whether or not the target operator is a jumping window operator (1104) representing the time (1104). If YES is determined in step 1104, the time information in which the target operator in the query execution unit 226 executed the previous process + the query execution tree graph analysis result management table 235 is displayed in the next firing time holding area 234. The registered value of the jumping window size is registered (1105).
When the step 1105 is completed, or when NO is determined in the step 1104, the next firing time calculation unit 223 determines whether the target operator is a streaming operator (IStream, DStream, IDStream) that causes a delay representing the time. Is determined (1106). If YES is determined in step 1106, the time information given to the tuple input to the target operator in the query execution unit 226 in the next firing time holding area 234 + the query execution tree graph analysis result management. The value of the delay size registered in Table 235 is registered (1107).
When the step 1107 is completed, or when NO is determined in the step 1106, the next firing time calculation unit 223 determines whether or not the target operator is a streaming operator (RStream) that outputs at regular intervals representing the time. Judge (1108). If YES is determined in step 1108, the time information in which the target operator in the query execution unit 226 executed the previous process + the query execution tree graph analysis result management table 235 is displayed in the next firing time holding area 234. The registered output interval value is registered (1109).
When the step 1109 is completed, or when NO is determined in the step 1108, the next firing time calculation unit 223 is an operator (Sum, Count, Average, Min, Max, Median) having a function of removing the ghost by the target operator. , Variable, Standard Deviation, Limit) (1110). If YES is determined in step 1110, the time information given to the tuple input to the target operator in the query execution unit 226 in the next firing time holding area 234 + the query execution tree graph analysis result management. Register the minimum time unit values registered in Table 235 (1111).
When the step 1111 is completed, or when NO is determined in the step 1110, the process of the step 804 is terminated (1112).
By the above processing, in the next firing time holding area 234, the sum of the time information registered in the query execution tree graph analysis result management table 235 is stored in the time information in which the operator in the query execution unit 226 executed the previous processing. Will be done. That is, the next firing time holding area 234 stores the next firing time, which is the time when each operator should start processing next time.
FIG. 12 is a flowchart showing the HBT insertion (or generation) process of step 805 shown in FIG.
First, the HBT generation unit 224 shown in FIG. 2 acquires the current system time held in the system time holding area 232 (1202). Next, the HBT generation unit 224 acquires the oldest next ignition system time held in the next ignition time holding area 234 (1203). Next, the HBT generation unit 224 compares the current system time with the next ignition system time, and the value of the current system time is equal to or greater than the value of the next ignition system time (current system time value> = next ignition. Determine if it is (system time) (1204).
If YES is determined in step 1204, the HBT generation unit 224 acquires the final data input time held in the HBT data processing time holding area 233, and the HBT generation unit 224 acquires the final data input time. And the next ignition system time acquired in step 1203 are compared, and it is determined whether or not the final data input time is smaller than the next ignition system time (final data input time <next ignition system time) (1206). ..
If YES is determined in step 1206, the value of the next ignition system time acquired in step 1202 as the final data input time held in the HBT generation data processing time holding area 233 by the HBT generation unit 224. (1207), and the HBT generation unit 224 transmits the HBT at the next ignition system time (1208).
When the step 1208 is completed, or when NO is determined in the step 1206, the next firing system time acquired in the next firing time 1202 held in the next firing time holding area 234 by the HBT generation unit 224 is set. Delete the value (1209).
When the step 1209 is completed, or when NO is determined in the step 1204, the process of the step 805 is terminated (1210).
Here, in step 1203, the oldest next ignition system time held in the next ignition time holding area 234 is acquired, but in step 1203, the time is equal to or less than the current time value and the latest next ignition time. And the time information equal to or less than the value of the next firing time may be deleted in step 1209.
By the above processing, when the condition that the current system time value is equal to or more than the value of the next ignition system time and the final data input time is smaller than the next ignition system time is satisfied, the HBT generator 224 sets the next ignition system. Send the HBT of the time. As a result, the execution tree 241 shown in Fig. 6<sub>1</sub>Operator can start processing without waiting for the arrival of stream data 21.
FIG. 13 shows a sequence diagram illustrating the next firing time calculation process and the HBT generation process.
In FIG. 13, the stream data receiving unit 221 and the HBT generating unit 224 shown in FIG. 6 are shown.<sub>1</sub>, Next firing time holding area 234<sub>1</sub>, HBT generation data processing time holding area 233<sub>1</sub>, HBT generator 224<sub>2</sub>, Next firing time holding area 234<sub>2</sub>, HBT generation data processing time holding area 233<sub>2</sub>, The window operator 611, the window operator 612, the next firing time calculation unit 223 shown in FIG. 2, and the system time holding area 232 will be described. Further, the flowcharts shown in FIGS. 8, 10, 11, and 12 will also be described. Here, HBT generator 224<sub>1</sub>, HBT generator 224<sub>2</sub>, Window operator 611, and window operator 612 are controlled by the query execution tree scheduler 222 shown in FIG.
First, the stream data 21 shown in FIG. 6 is connected to the stream data receiver 221.<sub>1</sub>Is entered (1301). Temperature stream data 21 as shown in Figures 3 and 4<sub>1</sub>The form of is the temperature sensor ID302, the device ID303, and the temperature 304. In the 1301, the value of the temperature sensor ID302 is "1001", the value of the device ID303 is "201", and the value of the temperature 304 is "23 ° C". To do. In 1301, the format is simplified and only "23 ° C" is shown.
Next, the stream data receiving unit 221 acquires the current time "9:00" held in the system holding area 232 (1302), stores "9:00" in the value of the system time stamp 305, and generates an HBT. Part 224<sub>1</sub>Send to (1303).
HBT generator 224<sub>1</sub>Is determined to be YES in step 1002 according to the processing of the flowchart shown in FIG. 10, and stream data 21 is determined in step 1003.<sub>1</sub>HBT generator 224 that accepted the tuple "9:00, 23 ° C"<sub>1</sub>Next firing time holding area 234 other than, that is, HBT generator 224<sub>2</sub>Next firing time holding area 234<sub>2</sub>Stream data received in 21<sub>1</sub>Register the input time information "9:00" of (1304). Next, HBT generator 224<sub>1</sub>Data processing time holding area for HBT generation 233<sub>1</sub>Stream data 21 that received the last data input time of<sub>1</sub>The input time information of is updated to the value of "9:00" (1305), and the tuple "9:00, 23 ° C" is transmitted to the window operator 611 (1306).
The window operator 611 executes the calculation process of the next firing time after executing the process unique to the window operator 611. According to the flowchart shown in FIG. 11, YES is determined in step 1102, and the next ignition time holding area 234<sub>1</sub>, Next firing time holding area 234<sub>2</sub>The time information "9:00" given to the tuple "9:00, 23 ° C" entered in the window operator 611 + registered in row 704 of the query execution tree graph analysis result management table 235 shown in Fig. 7. Register the value of the sliding window size "2 minutes" = "9:02" (1308, 1309). Here, the next firing time holding area 234<sub>2</sub>Will be added to, and the next ignition time holding area 234<sub>2</sub>The time information stored in is "9:00, 9:02". Next, in steps 1104, 1106, 1108, and 1110, all of them are determined to be NO, and the process ends. Further, the window operator 611 adds a plus flag indicating that the tuple has been processed, and transmits the tuple "+, 9:00, 23 ° C" to the selection operator 613 shown in FIG. 6 (the window). It is stored in the operator concatenation queue 236 that connects the operator 611 and the selected operator 613) (1311).
HBT generator 224<sub>2</sub>Acquires the current system time "9:00" held in the system time holding area 232 in step 1202 according to the flowchart shown in FIG. 12 (1321), and the next firing time holding area 234 in step 1203.<sub>2</sub>Gets the oldest next firing system time "9:00" held in (1322). In step 1204, YES is determined because the relationship of the current system time 9:00> = the next ignition system time 9:00 is established, and in step 1205, the data processing time holding area for HBT generation 233<sub>2</sub>Gets the last data entry time of "-(No data is stored because it has not been updated yet)" (1323). In step 1206, the last data input time- (no data is stored because it has not been updated yet) <the next ignition system time "9:00" is satisfied, so it is determined as YES, and the step 1206 Data processing time holding area for HBT generation in 1207 233<sub>2</sub>The last data input time held in is updated to the value of the next ignition system time "9:00" acquired in step 1322 (1324). Next, in step 1208, the HBT HBT, 9:00 of the next ignition system time 9:00 acquired in step 1322 is transmitted to the window operator 612 (1325), and the next ignition time is held in step 1209. Area 234<sub>2</sub>The value of the next ignition system time "9:00" acquired in step 1322 is deleted (1326). Since the window operator 612 is a Partitioned window, it does not perform any processing even if it receives the HBT "HBT, 9:00", and the HBT "HBT, 9:00" is shown in Fig. 6. It is transmitted to the indicated selected operator 614 (stored in the operator concatenation queue 236 that connects the window operator 612 and the selected operator 614) (1327).
HBT generator 224<sub>1</sub>Acquires the current system time "9:01" held in the system time holding area 232 in step 1202 according to the flowchart shown in FIG. 12 (1331), and obtains the next firing time holding area 234 in step 1203.<sub>1</sub>Gets the oldest next firing system time "9:02" held in (1332). In step 1204, NO is determined because the relationship of the current system time 9:01> = the next ignition system time 9:02 does not hold, and the process ends.
Similarly, HBT generator 224<sub>2</sub>Acquires the current system time "9:01" held in the system time holding area 232 in step 1202 according to the flowchart shown in FIG. 12 (1341), and obtains the next firing time holding area 234 in step 1203.<sub>1</sub>Gets the oldest next firing system time "9:02" held in (1342). In step 1204, NO is determined because the relationship of the current system time 9:01> = the next ignition system time 9:02 does not hold, and the process ends.
HBT generator 224<sub>1</sub>Acquires the current system time "9:02" held in the system time holding area 232 in step 1202 according to the flowchart shown in FIG. 12 (1351), and the next firing time holding area 234 in step 1203.<sub>1</sub>Gets the oldest next firing system time "9:02" held in (1352). In step 1204, YES is determined because the relationship of the current system time 9:02> = the next ignition system time 9:02 is established, and in step 1205, the data processing time holding area for HBT generation 233<sub>1</sub>Get the last data entry time "9:00" of (1353). In step 1206, YES is determined because the relationship of the final data input time 9:00 <the next ignition system time 9:02 is established, and in step 1207, the data processing time holding area for HBT generation 233<sub>1</sub>The last data input time held in is updated to the value of the next ignition system time "9:02" acquired in step 1352 (1354). Next, in step 1208, the HBT HBT, 9:02 of the next ignition system time 9:02 acquired in step 1352 is transmitted to the window operator 611 (1355), and the next ignition time is held in step 1209. Area 234<sub>1</sub>The value of the next firing system time "9:02" acquired in step 1352, which is held in, is deleted (1356). The window operator 611 transmits the HBT HBT, 9:02 to the select operator 613 shown in FIG. 6 (stored in the operator concatenation queue 236 connecting the window operator 611 and the select operator 613) (1357). Also, due to the HBT "HBT, 9:02", the tuple "+, 9:00, 23 ° C" transmitted in step 1311 is excluded from the processing target (because of the definition of Range 2 minute, the tuple at 9:00). Is excluded from the processing target at 9:02), so the tuple "-, 9:00, 23 ° C" with a minus flag (-) indicating that it was excluded from the processing target is sent to the selection operator 613 shown in Fig. 6. Send (stored in the operator concatenation queue 236 that connects the window operator 612 and the selected operator 613) (1358).
Similarly, HBT generator 224<sub>2</sub>Acquires the current system time "9:02" held in the system time holding area 232 in step 1202 according to the flowchart shown in FIG. 12 (1361), and the next firing time holding area 234 in step 1203.<sub>2</sub>Gets the oldest next firing system time "9:02" held in (1362). In step 1204, YES is determined because the relationship of the current system time 9:02> = the next ignition system time 9:02 is established, and in step 1205, the data processing time holding area for HBT generation 233<sub>2</sub>Get the last data entry time "9:00" of (1363). In step 1206, YES is determined because the relationship of the final data input time 9:00 <the next ignition system time 9:02 is established, and in step 1207, the data processing time holding area for HBT generation 233<sub>2</sub>The last data input time held in is updated to the value of the next ignition system time "9:02" acquired in step 1362 (1364). Next, in step 1208, the HBT HBT, 9:02 of the next ignition system time 9:02 acquired in step 1362 is transmitted to the window operator 612 (1365), and the next ignition time is held in step 1209. Area 234<sub>2</sub>The value of the next firing system time "9:02" acquired in step 1362 is deleted (1366). The window operator 612 transmits the HBT HBT, 9:02 to the select operator 614 shown in FIG. 6 (stored in the operator concatenation queue 236 connecting the window operator 612 and the select operator 614) (1367). At this time, even if the window operator 612 receives the HBT HBT, 9:02, the processing target does not change.
Next, the stream data 21 shown in FIG. 6 is connected to the stream data receiver 221.<sub>2</sub>Is entered (1371). Temperature stream data 21 as shown in Figures 3 and 4<sub>2</sub>The form of is humidity sensor ID 312, device ID 313, and humidity 314. In the above 1371, the value of the humidity sensor ID 312 is "2001", the value of the device ID 313 is "201", and the value of the humidity 304 is "67%". .. In 1371, the format is simplified and only "67%" is shown.
Next, the stream data receiving unit 221 acquires the current time "9:03" held in the system holding area 232 (1372), stores "9:03" in the value of the system time stamp 315, and generates HBT. Part 224<sub>2</sub>Send to (1373).
Hereinafter, the same process is repeated.
<Summary of the first embodiment> In the stream data processing method in which HBT indicating that the time has advanced is inserted (generated) in addition to the stream data by processing the data for which the processing target of the stream data is defined by the window, the received stream data The time information is held in the next firing time holding area as the firing time, which is the time when the HBT generation unit that generates the HBT inserts the HBT. Further, a processing module that generates the HBT is extracted from a query representing the processing method of the stream data at a time different from the time when the stream data is received, and the extracted processing module and the time information when the stream data is received are extracted. The ignition time is calculated from the above and held in the next ignition time holding area. Then, at the ignition time held in the next ignition time holding region, the HBT generation unit generates the HBT. It was shown that the time control information, which is the object of the present invention, can be inserted when necessary by the above processing.
Further, as shown above, in the process of the present invention, the time control information may be inserted when necessary, and the amount of the time control information can be reduced. However, since the time control information is inserted at the timing when the processing module is required, stream data processing with low latency can be realized.
The first embodiment of the present invention has been described above.
The present invention is not limited to the first embodiment shown above, and many modifications can be made within the scope of the gist thereof. In the following, it is possible to obtain the same or further effect by an embodiment different from the first embodiment, or to obtain a further effect by combining with the first embodiment. Will be explained.
For example, in the flowchart shown in FIG. 9, the next ignition time calculation process shown in FIG. 13, and the sequence diagram illustrating the process at the time of receiving the HBT generation process, the next ignition time is set for all the next ignition regions. It was stored. However, as shown in the flowchart shown in FIG. 14, by determining whether or not to store the next ignition time for the next ignition area, the next ignition time is stored in all the next ignition time holding areas. It does not have to be.
The flowchart shown in FIG. 14 will be described below.
The stream data processing unit 220 is the all HBT generation unit 224.<sub>n</sub>The process shown below is repeated for the target (1402).
If YES is determined in the previous process of step 1403, the target HBT generator 224<sub>n</sub>Determines whether there is a parent-child relationship with the target operator who is trying to store the next firing time and the graph structure of the execution tree (1403). The parent-child relationship is the execution tree 241 shown in FIG.<sub>1</sub>In, it is determined whether or not the two target operators are on the path through which the stream data 21 passes. When handling stream data with two or more inputs such as the join operator 615, the operator on the output side has a parent-child relationship with the operator of any input. For example, in FIG. 6, the aggregation operator 616 and the HBT generation unit 224<sub>1</sub>Is a parent-child relationship.
When the step 1404 is completed, or when NO is determined in the step 1403, the target HBT generator 224<sub>n</sub>Next firing time holding area 234<sub>n</sub>The next firing time is registered in (1404). All HBT generators 224<sub>n</sub>When the steps 1403 and 1404 are processed, the process ends (1405).
In the flowchart shown in FIG. 14, at the time of storing the next firing time, it is determined whether or not to store in the next firing time holding area, but in advance which next firing time holding area the target operator stores is determined. , It may be stored in a table and the table may be referred to.
<Second embodiment> The second embodiment of the present invention will be described below.
In the first embodiment, the stream data processing server 100 generates an HBT for advancing the time even during a period when no data is generated, and the HBT generation unit 224 generates the HBT when necessary, thereby generating the query execution unit. It was shown that the processing wait of 226 can be eliminated.
In the second embodiment, even when stream data processing is performed in the second computer by using the time information in the first computer in a plurality of computers, the time information in the first computer is transmitted to the second computer. This causes a problem similar to the problem of waiting for processing of the query execution unit 226 described above.
In the second embodiment, when performing stream data processing in the second computer using the time information in the first computer in a plurality of computers, the time information in the first computer is required for the second computer. It is characterized by transmitting time control information. In the second embodiment, the time control information is referred to as a system time stamp tuple (STT) to distinguish it from the HBT used inside the query execution unit 226. The STT has the same form as the HBT, and has an STT flag indicating that the STT is an STT and information on the time of occurrence. The STT may be in another form.
As shown in the first embodiment, the second computer calculates the next ignition time and transmits it to the first computer. The first computer transmits the STT to the second computer based on the next firing time. The query execution in the second computer may be performed by the processing method shown in the first embodiment or any other method. For example, the query may be processed without using HBT.
FIG. 15 is a block diagram showing a stream data processing server to which the second embodiment of the present invention is applied and a related system configuration.
The stream data processing server 100 is the same as the stream data processing server 100 shown in FIG. 1, and is composed of a CPU 11, a memory 12, a DISK 13, and an I / F 14, and the memory 12 includes an operating system (OS) 200 and commands. It is composed of an input unit 210 and a stream data processing unit 220. The operating system (OS) 200 and the command input unit 210 are the same as the operating system (OS) 200 and the command input unit 210 shown in FIG. The stream data processing unit 220A of the stream data processing server 100 deletes the system time holding area and the next firing time holding area from the stream data processing unit 220 of the first embodiment, and adds the STT receiving unit 1562.
The application operation server 1500 is composed of a CPU 1501, a memory 1502, a DISK1503, and an I / F 1504, and the memory 1502 is composed of an operating system (OS) 1510, a command input unit 1520, and a stream data generation application 1530. The stream data processing server 100 and the application operation server 1500 are connected to the network 112 shown in FIG. 1 via I / F14 and I / F1504.
The outline of the present invention will be described with reference to FIG. The application operation server 1500 generates stream data 21, adds time information held by the application operation server 1500, and transmits the stream data to the stream data processing server 100. The stream data processing server 100 reads the stream data 21 based on an inquiry commanded by the user 114 from the computer 115 or an inquiry input by an application executed on the computer 115, and is held by the application operation server 1500. After converting to meaningful information based on the time information, the output result 23 is output. Here, the significant information is, for example, the information obtained by converting the measured value of the temperature sensor node 101 in FIG. 1 into an average value for a certain period of time because the users 114 and 116 cannot understand the measured value data series as it is. ..
Hereinafter, the configuration of the stream data processing server 100 will be described in detail.
The stream data processing unit 220 includes a stream data receiving unit 1561, an STT receiving unit 1562, a stream data processing unit 1563, a next firing time calculation unit 1564, a query execution tree graph analysis unit 1565, a query execution time holding area 1571, and input stream data. It consists of a retention buffer 1572, a query execution tree analysis result management table 1573, an operator concatenation queue 1574, and an output result retention buffer 1575.
The input stream data holding buffer 1572 is the same as the input stream data holding buffer 231.
The stream data receiving unit 1561 acquires the data of the input stream data holding buffer 1572 and outputs the data to the stream data processing execution unit 1563. In this embodiment, since the processing is performed based on the time information given by the application operation server, the current time of the system held by the stream data processing server 100 is not used.
The stream data processing execution unit 1563 processes the stream data 21 based on the time information given by the application operation server. Any processing method may be used. For example, a processing unit such as the query execution unit 226, the query execution tree scheduler 222, and the HBT generation unit 224 shown in FIG. 2 may be combined. In the present embodiment, it is assumed that the operators such as the window operation, the selection operation, the projection operation, the join operation, and the aggregation operation as shown in the query execution unit 226 have a tree structure (execution tree). The stream data processing execution unit 1563 receives the data output from the stream data reception unit 1561, and stores the result processed by each operator of the execution tree in the output result holding buffer 1575. The intermediate result processed by each operator is stored in the operator concatenation queue 1574.
The operator concatenated queue 1574 and the output result holding buffer 1575 are the same as the operator concatenated queue 236 and the output result holding buffer 237. Further, the query execution tree analysis unit 1565 and the query execution tree analysis result management table 1573 are the same as the query execution tree analysis unit 225 and the query execution tree analysis result management table 235.
The next firing time calculation unit 1564 refers to the query execution tree analysis result management table 1573, and is based on the time information of the input stream data and the setting information of the query stored in the query execution tree analysis result management table 1573. The next firing time is calculated, and the next firing time is transmitted to the application operation server 1500 via the I / F 14 as the next firing time message. The next firing time calculation unit 1564 is called when the stream data processing execution unit 1563 executes the operator extracted by the query execution tree analysis unit 1565.
The STT receiving unit 1562 receives the STT transmitted from the STT generating unit 1542 of the application operation server 1500, which will be described later, via the I / F14, and holds the time information given to the STT in the inquiry execution time holding area 1571. To do. The inquiry execution time holding area 1571 is an area for holding time information transmitted from the application operation server 1500, and the stream data processing execution unit 1563 executes a query based on the time information.
Next, the configuration of the application operation server 1500 will be described in detail.
The command input unit 1520 receives a command commanded by the user 114 from the computer 115 or a command input by an application executed on the computer 115.
The stream data generation application 1530 is composed of a stream data generation unit 1531, a next firing time message receiving unit 1541, an STT generation unit 1542, a system time holding area 1551, and a next firing time holding area 1552.
The system time holding area 1551 is an area for holding the current time of the system. In the present embodiment, the current time of the system is the absolute time information (for example, the current time managed by the OS 1510) possessed by the application operation server 1500. Further, the current time of the system may be a value updated by the time information input from another computer.
The stream data generation unit 1531 generates the stream data 21, assigns the current time of the system held in the system time holding area 1551 to the stream data 21, and sends the stream data processing server 100 to the stream data processing server 100 via the I / F 1504. Send.
The next firing time message receiving unit 1541 receives the next firing time message transmitted from the next firing time calculation unit 1564 of the application operation server 1500 via the I / F 1504, and the time information given to the next firing time message. Is held in the next firing time holding area 1552. The next firing time holding area 1552 is an area for holding the next firing time given to the next firing time message received by the next firing time message receiving unit 1541.
The STT generation unit 1542 refers to the system time holding area 1551 and the next firing time holding area 1552, and the details will be described later, but the current time of the system held in the system time holding area 232 and the next firing time holding. From the next ignition time held in the area 234, the STT is transmitted to the stream data processing server 100 via the I / F 1504 at the ignition time.
Here, the stream data 21, the output result 23, the next firing time message, the STT, and the temporarily stored data held by the operator for processing can be in any data format such as a tuple format (record format), an XML format, or a CSV file. Good. An example of using the tuple format will be described below. Further, the stream data 21, the output result 23, the STT, and the temporarily stored data held by the operator for processing do not need to have a data entity, and a pointer pointing to a part or all of the data is a pointer to the data entity. It may be included.
Further, the application operation server 1500 may be a server that executes stream data processing. For example, the stream data generation unit 1531 may be the stream data processing unit 220 shown in FIG. Then, the stream data processing is divided into the first computer and the second computer for processing, and the second computer executes the stream data processing using the time information obtained by the stream data processing by the first computer. May be good.
FIG. 16 is a flowchart showing the overall processing of the stream data processing server 100 and the application operation server 1500.
First, the query execution tree analysis unit 1565 of the stream data processing server 100 shown in FIG. 15 extracts a firing operator from the query execution tree of the registered query and registers it in the query execution tree analysis result management table 1573 (1602). .. The detailed contents of the process of step 1602 are the same as the flowchart shown in FIG.
Next, the next firing time calculation unit 1564 of the stream data processing server 100 calculates the next firing time, and sends the next firing time message to the application operation server 1500 (1603). The details of the process of step 1603 will be described later with reference to FIG.
Next, the next firing time message receiving unit 1541 of the application operating server 1500 receives the next firing time message and holds it in the next firing time holding area 1552 of the application operating server 1500 (1604).
Next, the STT is transmitted to the stream data processing server 100 at the firing time held in the next firing time holding area 1552 by the STT generation unit 1542 of the application operation server 1500 (1605). The details of the process of step 1605 will be described later with reference to FIG.
Next, it is determined whether or not the system termination command has been accepted by the command input unit 210 of the stream data processing server 100 (1606). If NO is determined in step 1606, the process returns to step 1602, and if YES is determined in step 1606, the processing of the stream data processing server 100 is terminated (1607).
FIG. 17 is a flowchart showing the next firing time calculation process of step 1603 shown in FIG.
First, the next firing time calculation unit 1564 of the stream data processing server 100 determines whether or not the target operator is a sliding window operator (Range Window) representing time (1702). If YES is determined in step 1702, the next firing time calculation unit 1564 is assigned to the tuple input to the target operator in the stream data processing execution unit 1563 shown in FIG. 15 to the application operation server 1500. The value of the sliding window size registered in the time information + the query execution tree graph analysis result management table 1573 is transmitted as the next firing time message (1703).
When the step 1703 is completed, or when the determination is NO in the step 1702, the next firing time calculation unit 1564 determines whether or not the target operator is a jumping window operator representing the time (1704). .. If YES is determined in step 1704, the next firing time calculation unit 1564 executes the previous processing by the target operator in the stream data processing execution unit 1563 on the application operation server 1500. The time information + the inquiry. The value of the jumping window size registered in the execution tree graph analysis result management table 1573 is transmitted as the next firing time message (1705).
When the step 1705 is completed, or when the determination is NO in the step 1704, the next firing time calculation unit 1564 determines whether the target operator is a streaming operator (IStream, DStream, IDStream) that causes a delay representing the time. (1706). If YES is determined in step 1706, the next firing time calculation unit 1564 gives the time information given to the tuple input to the target operator in the stream data processing execution unit 1563 to the application operation server 1500. + The value of the delay size registered in the query execution tree graph analysis result management table 1573 is transmitted as the next firing time message (1707).
When the step 1707 is completed, or when the determination is NO in the step 1706, the next firing time calculation unit 1564 is whether or not the target operator is a streaming operator (RStream) that outputs at regular intervals representing the time. Is determined (1708). If YES is determined in step 1708, the next firing time calculation unit 1564 executes the previous processing by the target operator in the stream data processing execution unit 1563 on the application operation server 1500. The time information + the inquiry. The value of the output interval registered in the execution tree graph analysis result management table 1573 is transmitted as the next firing time message (1709).
When the step 1709 is completed, or when the determination is NO in the step 1708, the next firing time calculation unit 1564 determines that the target operator has a function of removing ghosts (Sum, Count, Average, Min, Max, Median, Variable, Standard Deviation, Limit) (1710). If YES is determined in step 1710, the next firing time calculation unit 1564 gives the time information given to the tuple input to the target operator in the stream data processing execution unit 1563 to the application operation server 1500. + The value in the minimum time unit registered in the query execution tree graph analysis result management table 1573 is sent as the next firing time message (1711).
When the step 1711 is completed, or when NO is determined in the step 1710, the process of the step 1603 is terminated (1712).
By the above processing, the next firing time of the firing operator is calculated from the sum of the time information given to the tuple input to the target operator and the time information registered in the query execution tree graph analysis result management table 1573. It is sent to the application operation server 1500 as the next firing time message.
FIG. 18 is a flowchart showing the processing of the STT generation unit 1542 of the application operation server 1500 performed in step 1605 shown in FIG.
First, the STT generator 1542 shown in FIG. 15 acquires the current system time held in the system time holding area 1551 (1802). Next, the STT generation unit 1542 acquires the oldest next ignition system time held in the next ignition time holding area 1551 (1803). Next, the STT generator 1542 compares the current system time with the next ignition system time, and the value of the current system time is equal to or greater than the value of the next ignition system time (current system time value> = next time). Ignition system time) is determined (1804).
If YES is determined in step 1804, the STT generation unit 1542 transmits the STT of the next ignition system time (1805), and the STT generation unit 1542 holds the STT in the next ignition time holding area 1551. Delete the value of the next firing system time obtained in step 1802 (1806).
When the step 1806 is completed, or when NO is determined in the step 1804, the process of the step 1605 is terminated (1807).
Here, in step 1803, the oldest next ignition system time held in the next ignition time holding area 1551 is acquired, but in step 1803, the value is equal to or less than the current time value and the latest next ignition is made. The time may be acquired and the time information equal to or less than the value of the next firing time may be deleted in step 1806.
The second embodiment of the present invention has been described above.
<Summary> The present invention is not limited to the first and second embodiments shown above, and a number of modifications can be made within the scope of the gist thereof. Further, an embodiment in which the first and second embodiments shown above are combined is also possible.
For example, in the above embodiment, the amount of data to be held in the next firing time holding area shown in FIGS. 2 and 15 is not limited, but it is appropriate when the upper limit of the memory is set and the upper limit is exceeded. Processing may be performed. For example, old data of time information may be deleted, input stream data may be temporarily stopped, or a part of input stream data may be deleted and processed (shredded).
Further, in the above embodiment, the next firing time calculation unit and the query execution tree analysis unit have been described as being in the stream data processing server 100, but they may be processed by a computer different from the stream data processing server 100. Absent.
Further, in the above embodiment, an example of processing time control information (HBT and STT) in the stream data processing server has been described, but it is shown in the above embodiment in a system other than the stream data processing server such as a database system. The time control information may be processed.
Further, in the above embodiment, the stream data processing server 100 and the application operation server 1500 have been described as an arbitrary computer system, but a part of the processing performed by the stream data processing server 100 and the application operation server 1500. Alternatively, the entire process may be performed by the storage device.
Further, in the above embodiment, an example in which the sensor base station 108 inputs temperature data and humidity data as stream data 21 to the stream data processing server 100 has been described, but the present invention is not limited to this. For example, a sensor net server that manages a large number of sensor nodes on behalf of the sensor base station 108 outputs the measured values of the sensor nodes as stream data 21, and the stream data processing server 100 outputs meaningful information that can be understood by the user 116. It may be converted into the including output result 23 and provided to the computer 117. Further, the data input to the stream data processing server 100 may be the tag information read by the RFID reader or the data input from the computer 113, which is an RFID middleware system that centrally manages RFID. Further, the data may be input from the stock information providing server 118. In addition, traffic information such as ETC system, IC card information such as automatic ticket gates and credit cards, financial information such as stock price information, manufacturing process management information, call information, system log, network access information, traceability individual item information, Surveillance video metadata, Web clickstream, etc. may be used.
As described above, the present invention can realize stream data processing with low latency while reducing the amount of time control information by inserting (or generating) the time control information when necessary. In particular, it can be applied to financial applications, traffic information systems, traceability systems, sensor monitoring systems, computer system management, etc., where the amount of stream data that needs to be processed in real time is enormous.
<figref num="1">It is a block diagram which shows an example of the computer system of this invention.</figref><figref num="2">It is a block diagram which shows the structure of the stream data processing server which showed the 1st Embodiment of this invention.</figref><figref num="3a">It is a figure which shows the 1st Embodiment of this invention, and schematically represented the example of the preferable data format of stream data 21.</figref><figref num="3b">It is a figure which shows the 1st Embodiment of this invention, and schematically represented the other example of the preferable data format of stream data 21.</figref><figref num="4">It is explanatory drawing which shows the 1st Embodiment of this invention and shows the description example of the suitable command at the time of registering and setting stream data 21 in a stream data processing server 100.</figref><figref num="5">It is explanatory drawing which shows the 1st Embodiment of this invention and shows the description example of a suitable command at the time of registering and setting a query registration command in a stream data processing server 100.</figref><figref num="6">It is explanatory drawing which shows the 1st Embodiment of this invention and showed an example of the inquiry execution part 226.</figref><figref num="7">It is the figure which showed the 1st Embodiment of this invention and showed the structural example of the query execution tree analysis result management table 235.</figref><figref num="8">It is a flowchart which shows the 1st Embodiment of this invention and showed the whole process.</figref><figref num="9">It is the flowchart which showed the 1st Embodiment of this invention and showed the processing procedure of the ignition operator extraction processing.</figref><figref num="10">It is a flowchart which shows the 1st Embodiment of this invention and showed the processing procedure which holds the stream data reception time information.</figref><figref num="11">It is a flowchart which shows the 1st Embodiment of this invention and showed the processing procedure which calculates the next firing time.</figref><figref num="12">It is a flowchart which shows the 1st Embodiment of this invention and showed the processing procedure which generates HBT.</figref><figref num="13">It is a sequence diagram which shows the 1st Embodiment of this invention and exemplifies the next firing time calculation process, and HBT generation process.</figref><figref num="14">It is a flowchart which shows the 1st Embodiment of this invention and showed the processing procedure which selects the next firing time holding area which stores the next firing time.</figref><figref num="15">It is a block diagram which shows the structure of the stream data processing server which showed the 2nd Embodiment of this invention.</figref><figref num="16">It is a flowchart which showed the whole processing of the 2nd Embodiment of this invention.</figref><figref num="17">It is a flowchart which shows the 2nd Embodiment of this invention and showed the processing procedure which transmits the next firing time.</figref><figref num="18">It is a flowchart which shows the 2nd Embodiment of this invention and showed the processing procedure which transmits STT.</figref>
Code description
100 stream data processing server 11 CPU 12 MEMORY 13 DISK 14 I / F 21 Stream data 22 commands 23 Output result 101 temperature sensor node 102 Humidity sensor node 103 RFID tag 104 mobile phone 105 network 106 network 107 Network 108 sensor base station 109 cradle 110 RFID reader 111 Mobile phone base station 112 network 113 Relay calculator 114 users 115 calculator 116 users 117 Calculator 118 Stock Offering Server 200 operating system (OS) 210 Command input section 220 Stream data processing unit 221 Stream data receiver 222 Query Execution Tree Scheduler 223 Next firing time calculation unit 224 HBT generator 225 Query Execution Tree Analysis Department 226 Query execution part 231 Input stream data retention buffer 232 System time holding area 233 Data processing time holding area for HBT generation 234 Next firing time holding area 235 Query execution tree analysis result management table 236 Operator concatenation queue 237 Output result retention buffer
4 members in 2 offices
Priority claims2
| Document | Office | Kind | Date |
|---|---|---|---|
| 2008276685 | Japan | A | |
| JP20080276685 | – | – | – |
Members4
| Document | Office | Kind | |
|---|---|---|---|
| US2010106853A1 | United States of America | A1 | |
| JP2010108044A | Japan | A | |
| US8095690B2 | United States of America | B2 | |
| JP5154366B2This record | Japan | B2 |
10 legal events, as the office reported them to INPADOC
Over the term
Point at a mark for the eventEvents
| Event | Code | |
|---|---|---|
| Cancellation because of no payment of annual feesLAPS | LAPS | |
| Renewal fee payment (event date is renewal date of database)FPAY | FPAY | |
| Certificate of patent or registration of utility modelJAPANESE INTERMEDIATE CODE: R150R150 | R150 | |
| Certificate of patent or registration of utility modelJAPANESE INTERMEDIATE CODE: R150R150 | R150 | |
| First payment of annual fees (during grant procedure)JAPANESE INTERMEDIATE CODE: A61A61 | A61 | |
| Written decision to grant a patent or to grant a registration (utility model)JAPANESE INTERMEDIATE CODE: A01A01 | A01 | |
| Report on retrievalJAPANESE INTERMEDIATE CODE: A971007A977 | A977 | |
| Decision of grant or rejection writtenTRDD | TRDD | |
| Written amendmentJAPANESE INTERMEDIATE CODE: A523A521 | A521 | |
| Written request for application examinationJAPANESE INTERMEDIATE CODE: A621A621 | A621 |
Numbers
- Publication
- 5154366
- Publication, DOCDB
- 5154366
- Publication, EPODOC
- JP5154366B
- Application
- 276685
- Application, DOCDB
- 2008276685
- Application, EPODOC
- JP20080276685
Titles2
- English
- Stream data processing program and computer system
- Japanese
- ストリームデータ処理プログラム及び計算機システム
Classification
- CPC, 3
- G06Q40/04
- G06Q10/06
- G06Q10/109
- IPC, 1
- G06F17 30